From 39e81a18fd7a8c58b9d86ad32ae150540c2f55b9 Mon Sep 17 00:00:00 2001 From: Jaremy Creechley Date: Mon, 21 Sep 2026 18:21:33 +0300 Subject: [PATCH 1/2] fix managed RChan ownership --- CHANGES.md | 8 + figdraw.nimble | 2 +- src/figdraw/common/imgutils.nim | 6 +- src/figdraw/common/rchannels.nim | 362 +++++++++++++++++-------------- tests/config.nims | 1 + tests/trchannels.nim | 356 ++++++++++++++++++++++++++++++ 6 files changed, 567 insertions(+), 168 deletions(-) create mode 100644 tests/trchannels.nim diff --git a/CHANGES.md b/CHANGES.md index 942cb9b2..279697b8 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -1,5 +1,13 @@ # Changes +## 0.41.0 + +- Fix managed-value ownership in `RChan` by storing queued values in typed + storage, including overwritten ring-buffer entries and unread messages. +- Release image-message subscription channels and run payload destructors + outside the channel lock, including the `tryTake` commit-race rollback. +- Add Linux/macOS RSS regressions covering both `tryRecv` and blocking `recv`. + ## 0.40.2 - Export producer-owned `RenderFragments`, fragment handles, cursors, and their diff --git a/figdraw.nimble b/figdraw.nimble index 979cfd42..8469f501 100644 --- a/figdraw.nimble +++ b/figdraw.nimble @@ -1,4 +1,4 @@ -version = "0.40.2" +version = "0.41.0" author = "Jaremy Creechley" description = "UI Engine for Nim" license = "MIT" diff --git a/src/figdraw/common/imgutils.nim b/src/figdraw/common/imgutils.nim index d742f4d5..def78b27 100644 --- a/src/figdraw/common/imgutils.nim +++ b/src/figdraw/common/imgutils.nim @@ -108,10 +108,14 @@ imageCachedLock.initLock() ownerTokenLock.initLock() imageSubscriberLock.initLock() -proc `=destroy`(subscription: ImageMessageSubscriptionHandle) = +proc `=destroy`(subscription: var ImageMessageSubscriptionHandle) = if subscription.id != 0'u64: withLock imageSubscriberLock: imageSubscribers.del(subscription.id) + # This custom destructor replaces the compiler-generated field cleanup. + # Release the subscription's channel owner after removing the table owner; + # queued ImageMsg values then get destroyed by RChan's final-owner cleanup. + `=destroy`(subscription.inbox) proc `==`*(a, b: ImageId): bool {.borrow.} proc `==`*(a, b: OwnerToken): bool {.borrow.} diff --git a/src/figdraw/common/rchannels.nim b/src/figdraw/common/rchannels.nim index c902adc0..d08632ee 100644 --- a/src/figdraw/common/rchannels.nim +++ b/src/figdraw/common/rchannels.nim @@ -26,7 +26,7 @@ ## ## The `RChan` type represents a generic fixed-size channel object that internally manages ## the underlying resources and synchronization. It has to be initialized using -## the `newChan` proc. Sending and receiving operations are provided by the +## the `newRChan` proc. Sending and receiving operations are provided by the ## blocking `send` and `recv` procs, and non-blocking `trySend` and `tryRecv` ## procs. For ring buffer behavior, use the `push` proc rather than `send`. ## Send operations add messages to the channel, receiving operations remove them, @@ -46,7 +46,7 @@ runnableExamples("--threads:on --gc:orc"): # Channels are generic, and they include support for passing objects between # threads. # Note that isolated data passed through channels is moved around. - var RChan = newChan[string]() + var RChan = newRChan[string]() block example_blocking: # This proc will be run in another thread. @@ -100,7 +100,7 @@ runnableExamples("--threads:on --gc:orc"): assert messages.len >= 2 block example_non_blocking_overwrite: - var chanRingBuffer = newChan[string](elements = 1) + var chanRingBuffer = newRChan[string](elements = 1) chanRingBuffer.push("Hello") chanRingBuffer.push("World") var msg = "" @@ -113,187 +113,213 @@ when not (defined(gcArc) or defined(gcOrc) or defined(gcAtomicArc) or defined(ni "This module requires one of --mm:arc / --mm:atomicArc / --mm:orc compilation flags" .} -import std/[locks, isolation, atomics] +import std/[atomics, deques, isolation, locks] # Channel # ------------------------------------------------------------------------------ type - ChannelRaw = ptr ChannelObj - ChannelObj = object + RChanItemData[T] = object + value: T + + RChanItem[T] = ptr RChanItemData[T] + + RChanData[T] = object lock: Lock spaceAvailableCV, dataAvailableCV: Cond - slots: int ## Number of item slots in the buffer - head: Atomic[int] ## Write/enqueue/send index - tail: Atomic[int] ## Read/dequeue/receive index + items: Deque[RChanItem[T]] + capacity: int + pendingSends: int atomicCounter: Atomic[int] - buffer: ptr UncheckedArray[byte] - -# ------------------------------------------------------------------------------ - -proc getTail(RChan: ChannelRaw, order: MemoryOrder = moRelaxed): int {.inline.} = - RChan.tail.load(order) - -proc getHead(RChan: ChannelRaw, order: MemoryOrder = moRelaxed): int {.inline.} = - RChan.head.load(order) - -proc setTail(RChan: ChannelRaw, value: int, order: MemoryOrder = moRelaxed) {.inline.} = - RChan.tail.store(value, order) - -proc setHead(RChan: ChannelRaw, value: int, order: MemoryOrder = moRelaxed) {.inline.} = - RChan.head.store(value, order) - -proc setAtomicCounter( - RChan: ChannelRaw, value: int, order: MemoryOrder = moRelaxed -) {.inline.} = - RChan.atomicCounter.store(value, order) - -proc numItems(RChan: ChannelRaw): int {.inline.} = - result = RChan.getHead() - RChan.getTail() - if result < 0: - inc(result, 2 * RChan.slots) - - assert result <= RChan.slots - -template isFull(RChan: ChannelRaw): bool = - abs(RChan.getHead() - RChan.getTail()) == RChan.slots - -template isEmpty(RChan: ChannelRaw): bool = - RChan.getHead() == RChan.getTail() -# Channels memory ops -# ------------------------------------------------------------------------------ - -proc allocChannel(size, n: int): ChannelRaw = - result = cast[ChannelRaw](allocShared(sizeof(ChannelObj))) - - # To buffer n items, we allocate for n - result.buffer = cast[ptr UncheckedArray[byte]](allocShared(n * size)) - - initLock(result.lock) - initCond(result.spaceAvailableCV) - initCond(result.dataAvailableCV) - - result.slots = n - result.setHead(0) - result.setTail(0) - result.setAtomicCounter(0) - -proc freeChannel(RChan: ChannelRaw) = - if RChan.isNil: + RChan*[T] = object ## Typed channel + d: ptr RChanData[T] + +when defined(figdrawRChanTests): + var rchanLiveItems*: Atomic[int] + +# The old implementation stored values in an untyped byte buffer and used +# copyMem to move them. That is not valid for managed values: the compiler never +# sees the slot's references, so overwriting or receiving a value cannot destroy +# the old owner. Keep the queue's entries as pointers to individually owned, +# typed values. Pointer operations are safe under the lock; each pointed-to +# value is constructed and destroyed outside the lock so user destructors can +# call back into the channel. + +proc allocItem[T](value: var Isolated[T]): RChanItem[T] = + result = cast[RChanItem[T]](allocShared0(sizeof(RChanItemData[T]))) + when defined(figdrawRChanTests): + discard rchanLiveItems.fetchAdd(1, moRelaxed) + try: + result[].value = extract(value) + except: + try: + `=destroy`(result[].value) + finally: + deallocShared(result) + when defined(figdrawRChanTests): + discard rchanLiveItems.fetchSub(1, moRelaxed) + raise + +proc freeItem[T](item: RChanItem[T]) = + if item.isNil: + return + try: + `=destroy`(item[].value) + finally: + deallocShared(item) + when defined(figdrawRChanTests): + discard rchanLiveItems.fetchSub(1, moRelaxed) + +proc allocChannel[T](n: Positive): ptr RChanData[T] = + result = cast[ptr RChanData[T]](allocShared0(sizeof(RChanData[T]))) + result[].items = initDeque[RChanItem[T]]() + result[].capacity = n + result[].atomicCounter.store(1, moRelaxed) + initLock(result[].lock) + initCond(result[].spaceAvailableCV) + initCond(result[].dataAvailableCV) + +proc freeChannel[T](channel: ptr RChanData[T]) = + if channel.isNil: return - if not RChan.buffer.isNil: - deallocShared(RChan.buffer) - - deinitLock(RChan.lock) - deinitCond(RChan.spaceAvailableCV) - deinitCond(RChan.dataAvailableCV) - - deallocShared(RChan) - -# MPMC Channels (Multi-Producer Multi-Consumer) -# ------------------------------------------------------------------------------ - -template incrWriteIndex(RChan: ChannelRaw) = - atomicInc(RChan.head) - if RChan.getHead() == 2 * RChan.slots: - RChan.setHead(0) - -template incrReadIndex(RChan: ChannelRaw) = - atomicInc(RChan.tail) - if RChan.getTail() == 2 * RChan.slots: - RChan.setTail(0) - -proc channelSend( - RChan: ChannelRaw, data: pointer, size: int, blocking: static bool, overwrite: bool + # RChanData lives in shared raw storage, so its managed queue field and every + # pointed-to managed value need explicit destruction before the allocation is + # released. No channel lock is held while payload destructors run. + # Detach the queue first: a payload destructor is allowed to re-enter the + # channel, and must not be able to receive an item that teardown is already + # destroying. + var items = channel[].items + channel[].items = initDeque[RChanItem[T]]() + for item in items.items: + freeItem(item) + deinitCond(channel[].spaceAvailableCV) + deinitCond(channel[].dataAvailableCV) + deinitLock(channel[].lock) + deallocShared(channel) + +proc channelSend[T]( + channel: ptr RChanData[T], + value: var Isolated[T], + blocking: static bool, + overwrite: static bool, ): bool = - assert not RChan.isNil - assert not data.isNil - - when not blocking: - if RChan.isFull() and not overwrite: - return false - - acquire(RChan.lock) - - # check for when another thread was faster to fill - when blocking: - if RChan.isFull(): - if overwrite: - incrReadIndex(RChan) - else: - while RChan.isFull(): - wait(RChan.spaceAvailableCV, RChan.lock) + assert not channel.isNil + + when overwrite: + let overwriteItem = allocItem(value) + var dropped: RChanItem[T] + var overwriteEnqueued = false + try: + acquire(channel[].lock) + try: + if channel[].items.len == channel[].capacity: + # Move the old owner out while holding the lock, but free its value + # after the lock is released. User destructors may call back into the + # channel. + dropped = channel[].items.popFirst() + channel[].items.addLast(overwriteItem) + overwriteEnqueued = true + signal(channel[].dataAvailableCV) + finally: + release(channel[].lock) + except: + if not overwriteEnqueued: + freeItem(overwriteItem) + freeItem(dropped) + raise + freeItem(dropped) + return true else: - if RChan.isFull(): - release(RChan.lock) - return false - - assert not RChan.isFull() - - let writeIdx = - if RChan.getHead() < RChan.slots: - RChan.getHead() + acquire(channel[].lock) + if blocking: + while channel[].items.len + channel[].pendingSends >= channel[].capacity: + wait(channel[].spaceAvailableCV, channel[].lock) else: - RChan.getHead() - RChan.slots - - copyMem(RChan.buffer[writeIdx * size].addr, data, size) - - incrWriteIndex(RChan) - - signal(RChan.dataAvailableCV) - release(RChan.lock) - result = true - -proc channelReceive( - RChan: ChannelRaw, data: pointer, size: int, blocking: static bool -): bool = - assert not RChan.isNil - assert not data.isNil - - when not blocking: - if RChan.isEmpty(): + if channel[].items.len + channel[].pendingSends >= channel[].capacity: + release(channel[].lock) + return false + inc channel[].pendingSends + + release(channel[].lock) + var item: RChanItem[T] + try: + item = allocItem(value) + except: + acquire(channel[].lock) + dec channel[].pendingSends + signal(channel[].spaceAvailableCV) + release(channel[].lock) + raise + + var commitRejected = false + var sendEnqueued = false + try: + acquire(channel[].lock) + try: + dec channel[].pendingSends + if not blocking: + if channel[].items.len == channel[].capacity: + commitRejected = true + else: + while channel[].items.len == channel[].capacity: + wait(channel[].spaceAvailableCV, channel[].lock) + if not commitRejected: + channel[].items.addLast(item) + sendEnqueued = true + signal(channel[].dataAvailableCV) + finally: + release(channel[].lock) + except: + if not sendEnqueued: + freeItem(item) + raise + + if commitRejected: + try: + value = unsafeIsolate(move item[].value) + finally: + freeItem(item) return false + result = true - acquire(RChan.lock) +proc channelReceive[T]( + channel: ptr RChanData[T], value: var T, blocking: static bool +): bool = + assert not channel.isNil - # check for when another thread was faster to empty + var received: RChanItem[T] + acquire(channel[].lock) when blocking: - while RChan.isEmpty(): - wait(RChan.dataAvailableCV, RChan.lock) + while channel[].items.len == 0: + wait(channel[].dataAvailableCV, channel[].lock) else: - if RChan.isEmpty(): - release(RChan.lock) + if channel[].items.len == 0: + release(channel[].lock) return false - assert not RChan.isEmpty() - - let readIdx = - if RChan.getTail() < RChan.slots: - RChan.getTail() - else: - RChan.getTail() - RChan.slots - - copyMem(data, RChan.buffer[readIdx * size].addr, size) - - incrReadIndex(RChan) - - signal(RChan.spaceAvailableCV) - release(RChan.lock) + # Move the owner out while holding the lock, but assign it to the caller's + # destination after unlocking. That assignment destroys any value already + # owned by the destination; T's destructor could call back into the channel. + received = channel[].items.popFirst() + signal(channel[].spaceAvailableCV) + release(channel[].lock) + try: + value = move received[].value + finally: + freeItem(received) result = true # Public API # ------------------------------------------------------------------------------ -type RChan*[T] = object ## Typed channel - d: ChannelRaw - template frees(c) = if c.d != nil: # this `fetchSub` returns current val then subs - # so count == 0 means we're the last - if c.d.atomicCounter.fetchSub(1, moAcquireRelease) == 0: + # fetchSub returns the count before decrementing, so one means this is the + # final owner. + if c.d.atomicCounter.fetchSub(1, moAcquireRelease) == 1: freeChannel(c.d) when defined(nimAllowNonVarDestructor): @@ -332,7 +358,7 @@ proc trySend*[T](c: RChan[T], src: sink Isolated[T]): bool {.inline.} = ## ## Returns `false` if the message was not sent because the number of pending ## messages in the channel exceeded its capacity. - result = channelSend(c.d, src.addr, sizeof(T), false, false) + result = channelSend(c.d, src, false, false) if result: wasMoved(src) @@ -360,7 +386,7 @@ proc tryTake*[T](c: RChan[T], src: var Isolated[T]): bool {.inline.} = ## ## Returns `false` if the message was not sent because the number of pending ## messages in the channel exceeded its capacity. - result = channelSend(c.d, src.addr, sizeof(T), false, false) + result = channelSend(c.d, src, false, false) if result: wasMoved(src) @@ -374,8 +400,8 @@ proc tryRecv*[T](c: RChan[T], dst: var T): bool {.inline.} = ## backoff strategy to reduce contention and improve the success rate of ## operations. ## - ## Returns `false` and does not change `dist` if no message was received. - channelReceive(c.d, dst.addr, sizeof(T), false) + ## Returns `false` and does not change `dst` if no message was received. + channelReceive(c.d, dst, false) proc send*[T](c: RChan[T], src: sink Isolated[T]) {.inline.} = ## Sends the message `src` to the channel `c`. @@ -387,7 +413,7 @@ proc send*[T](c: RChan[T], src: sink Isolated[T]) {.inline.} = ## messages from the channel are removed. when defined(gcOrc) and defined(nimSafeOrcSend): GC_runOrc() - discard channelSend(c.d, src.addr, sizeof(T), true, false) + discard channelSend(c.d, src, true, false) wasMoved(src) template send*[T](c: RChan[T], src: T) = @@ -402,7 +428,7 @@ proc push*[T](c: RChan[T], src: sink Isolated[T]) {.inline.} = ## The memory of `src` is moved, not copied. when defined(gcOrc) and defined(nimSafeOrcSend): GC_runOrc() - discard channelSend(c.d, src.addr, sizeof(T), true, overwrite = true) + discard channelSend(c.d, src, false, true) wasMoved(src) template push*[T](c: RChan[T], src: T) = @@ -417,21 +443,25 @@ proc recv*[T](c: RChan[T], dst: var T) {.inline.} = ## ## If the channel does not contain any messages this will block the thread until ## a message get sent to the channel. - discard channelReceive(c.d, dst.addr, sizeof(T), true) + discard channelReceive(c.d, dst, true) proc recv*[T](c: RChan[T]): T {.inline.} = ## Receives a message from the channel. ## A version of `recv`_ that returns the message. - discard channelReceive(c.d, result.addr, sizeof(T), true) + discard channelReceive(c.d, result, true) proc recvIso*[T](c: RChan[T]): Isolated[T] {.inline.} = ## Receives a message from the channel. ## A version of `recv`_ that returns the message and isolates it. - discard channelReceive(c.d, result.addr, sizeof(T), true) + var value: T + discard channelReceive(c.d, value, true) + result = unsafeIsolate(move value) proc peek*[T](c: RChan[T]): int {.inline.} = ## Returns an estimation of the current number of messages held by the channel. - numItems(c.d) + acquire(c.d[].lock) + result = c.d[].items.len + release(c.d[].lock) proc newRChan*[T](elements: Positive = 30): RChan[T] = ## An initialization procedure, necessary for acquiring resources and @@ -439,4 +469,4 @@ proc newRChan*[T](elements: Positive = 30): RChan[T] = ## ## `elements` is the capacity of the channel and thus how many messages it can hold ## before it refuses to accept any further messages. - result = RChan[T](d: allocChannel(sizeof(T), elements)) + result = RChan[T](d: allocChannel[T](elements)) diff --git a/tests/config.nims b/tests/config.nims index 5fa5d1c5..3522bf07 100644 --- a/tests/config.nims +++ b/tests/config.nims @@ -1,2 +1,3 @@ --path:"../src" --define:"figdraw.vulkanReadback" +--define:"figdrawRChanTests" diff --git a/tests/trchannels.nim b/tests/trchannels.nim new file mode 100644 index 00000000..7e505791 --- /dev/null +++ b/tests/trchannels.nim @@ -0,0 +1,356 @@ +import std/[atomics, isolation, os, unittest] +import figdraw/common/rchannels + +var destroyed: Atomic[int] + +type Payload = object + id: int + +proc `=destroy`(payload: Payload) = + if payload.id != 0: + discard destroyed.fetchAdd(1) + +type ReentrantPayload = object + id: int + +var + reentrantChannel: RChan[ReentrantPayload] + reentrantAlias: RChan[ReentrantPayload] + reentrantReady: bool + reentrantReceive: bool + reentrantCallbackDone: bool + +proc `=destroy`(payload: ReentrantPayload) = + if reentrantReady and not reentrantCallbackDone: + reentrantCallbackDone = true + if reentrantReceive: + var ignored: ReentrantPayload + discard reentrantAlias.tryRecv(ignored) + else: + discard reentrantChannel.peek() + +type BlockingPayload = object + value: int + +type RaisingPayload = object + id: int + +var raiseOnSink: bool + +proc `=sink`(dest: var RaisingPayload, src: RaisingPayload) = + if raiseOnSink: + raise newException(ValueError, "test sink failure") + dest.id = src.id + +var + blockingSinkStarted: Atomic[bool] + allowBlockingSink: Atomic[bool] + blockingSinkRaise: bool + raceInitialDone: Atomic[bool] + allowRaceRetry: Atomic[bool] + raceResult: Atomic[bool] + raceRetryResult: Atomic[bool] + raceSource: Atomic[int] + +proc `=sink`(dest: var BlockingPayload, src: BlockingPayload) = + if blockingSinkRaise: + raise newException(ValueError, "rollback sink failure") + dest.value = src.value + if not blockingSinkStarted.load(moAcquire): + blockingSinkStarted.store(true, moRelease) + while not allowBlockingSink.load(moAcquire): + sleep(1) + +var raceChannel: RChan[BlockingPayload] + +proc tryTakeDuringPush() {.thread.} = + {.cast(gcsafe).}: + var source = isolate(BlockingPayload(value: 42)) + let succeeded = raceChannel.tryTake(source) + raceResult.store(succeeded, moRelease) + if not succeeded: + var restored = extract(source) + raceSource.store(restored.value, moRelease) + source = isolate(move restored) + raceInitialDone.store(true, moRelease) + if not succeeded: + while not allowRaceRetry.load(moAcquire): + sleep(1) + raceRetryResult.store(raceChannel.tryTake(source), moRelease) + +proc tryTakeWithRaisingRollback() {.thread.} = + {.cast(gcsafe).}: + var source = isolate(BlockingPayload(value: 42)) + try: + discard raceChannel.tryTake(source) + except ValueError: + raceRetryResult.store(true, moRelease) + +const RaceWaitLimit = 5000 + +proc waitForFlag(flag: var Atomic[bool]): bool = + for _ in 0 ..< RaceWaitLimit: + if flag.load(moAcquire): + return true + sleep(1) + +suite "RChan managed payload ownership": + test "tryRecv destroys received payloads exactly once": + const iterations = 100 + destroyed.store(0) + block: + let channel = newRChan[Payload](1) + var value: Payload + for id in 1 .. iterations: + channel.send(isolate(Payload(id: id))) + check channel.tryRecv(value) + check value.id == id + check destroyed.load() == iterations + + test "recv destroys received payloads exactly once": + const iterations = 100 + destroyed.store(0) + block: + let channel = newRChan[Payload](1) + for id in 1 .. iterations: + channel.send(isolate(Payload(id: id))) + let value = channel.recv() + check value.id == id + check destroyed.load() == iterations + + test "recv into a destination destroys the previous value": + const iterations = 100 + destroyed.store(0) + block: + let channel = newRChan[Payload](1) + var value: Payload + for id in 1 .. iterations: + channel.send(isolate(Payload(id: id))) + channel.recv(value) + check value.id == id + check destroyed.load() == iterations + + test "push destroys overwritten payloads": + destroyed.store(0) + block: + let channel = newRChan[Payload](1) + channel.push(isolate(Payload(id: 1))) + channel.push(isolate(Payload(id: 2))) + check channel.recv().id == 2 + check destroyed.load() == 2 + + test "the final shared owner destroys unread payloads": + destroyed.store(0) + block: + var retained: RChan[Payload] + block: + let channel = newRChan[Payload](1) + retained = channel + channel.send(isolate(Payload(id: 1))) + check destroyed.load() == 0 + check destroyed.load() == 1 + + test "payload destruction can re-enter the channel": + reentrantReady = false + reentrantReceive = false + reentrantCallbackDone = false + block: + let channel = newRChan[ReentrantPayload](1) + reentrantChannel = channel + reentrantReady = true + channel.send(isolate(ReentrantPayload(id: 1))) + var value: ReentrantPayload + check channel.tryRecv(value) + reentrantReady = false + reset(reentrantChannel) + + test "final-owner cleanup detaches queued items before destruction": + reentrantReady = true + reentrantReceive = true + reentrantCallbackDone = false + block: + let channel = newRChan[ReentrantPayload](1) + copyMem(addr reentrantAlias, unsafeAddr channel, sizeof(reentrantAlias)) + channel.send(isolate(ReentrantPayload(id: 1))) + wasMoved(reentrantAlias) + reentrantReady = false + reentrantReceive = false + + test "tryTake preserves its source when push wins the commit race": + blockingSinkStarted.store(false) + allowBlockingSink.store(false) + raceInitialDone.store(false) + allowRaceRetry.store(false) + raceResult.store(true) + raceRetryResult.store(false) + raceSource.store(0) + raceChannel = newRChan[BlockingPayload](1) + + var worker: Thread[void] + createThread(worker, tryTakeDuringPush) + let sinkStarted = waitForFlag(blockingSinkStarted) + if sinkStarted: + raceChannel.push(BlockingPayload(value: 99)) + allowBlockingSink.store(true, moRelease) + let initialDone = waitForFlag(raceInitialDone) + if initialDone and not raceResult.load(moAcquire): + var queued: BlockingPayload + check raceChannel.tryRecv(queued) + check queued.value == 99 + allowRaceRetry.store(true, moRelease) + joinThread(worker) + + check sinkStarted + check initialDone + if sinkStarted and initialDone: + check not raceResult.load(moAcquire) + check raceSource.load(moAcquire) == 42 + check raceRetryResult.load(moAcquire) + var retried: BlockingPayload + check raceChannel.tryRecv(retried) + check retried.value == 42 + var leftover: BlockingPayload + discard raceChannel.tryRecv(leftover) + raceChannel = RChan[BlockingPayload]() + + test "raising ownership hooks do not strand channel capacity": + let channel = newRChan[RaisingPayload](1) + var source = isolate(RaisingPayload(id: 1)) + raiseOnSink = true + expect ValueError: + discard channel.tryTake(source) + raiseOnSink = false + + check channel.peek() == 0 + channel.send(isolate(RaisingPayload(id: 2))) + check channel.recv().id == 2 + + channel.send(isolate(RaisingPayload(id: 3))) + var destination = RaisingPayload(id: 0) + raiseOnSink = true + expect ValueError: + channel.recv(destination) + raiseOnSink = false + + check channel.peek() == 0 + channel.send(isolate(RaisingPayload(id: 4))) + check channel.recv().id == 4 + + test "raising rollback hooks still free the rejected item": + let liveItemsBefore = rchanLiveItems.load(moAcquire) + blockingSinkStarted.store(false) + allowBlockingSink.store(false) + blockingSinkRaise = false + raceRetryResult.store(false) + raceChannel = newRChan[BlockingPayload](1) + + var worker: Thread[void] + createThread(worker, tryTakeWithRaisingRollback) + let sinkStarted = waitForFlag(blockingSinkStarted) + if sinkStarted: + raceChannel.push(BlockingPayload(value: 99)) + blockingSinkRaise = true + allowBlockingSink.store(true, moRelease) + joinThread(worker) + blockingSinkRaise = false + + check sinkStarted + check raceRetryResult.load(moAcquire) + var queued: BlockingPayload + check raceChannel.tryRecv(queued) + check queued.value == 99 + raceChannel = RChan[BlockingPayload]() + check rchanLiveItems.load(moAcquire) == liveItemsBefore + +when defined(linux): + import std/strutils + + proc rssKb(): int64 = + for line in lines("/proc/self/status"): + if line.startsWith("VmRSS:"): + let fields = line.splitWhitespace() + if fields.len >= 2: + return parseInt(fields[1]).int64 + raise newException(IOError, "unable to read VmRSS from /proc/self/status") + +elif defined(macosx): + type + MachPort = uint32 + MachMsgTypeNumber = uint32 + TimeValue = object + seconds: int32 + microseconds: int32 + + MachTaskBasicInfo = object + virtualSize: uint64 + residentSize: uint64 + residentSizeMax: uint64 + userTime: TimeValue + systemTime: TimeValue + policy: int32 + suspendCount: int32 + + const machTaskBasicInfoFlavor = 20 + + var machTaskSelf {.importc: "mach_task_self_", header: "".}: + MachPort + + proc taskInfo( + task: MachPort, + flavor: cint, + taskInfoOut: ptr MachTaskBasicInfo, + taskInfoOutCount: ptr MachMsgTypeNumber, + ): cint {.importc: "task_info", header: "".} + + proc rssKb(): int64 = + var info: MachTaskBasicInfo + var count = MachMsgTypeNumber(sizeof(MachTaskBasicInfo) div sizeof(cuint)) + let status = taskInfo(machTaskSelf, machTaskBasicInfoFlavor, addr info, addr count) + doAssert status == 0, "task_info failed with kern_return_t " & $status + int64(info.residentSize div 1024) + +when defined(linux) or defined(macosx): + const + RssSampleCount = 6 + RssIterationsPerSample = 64 + RssPayloadSize = 512 * 1024 + RssMaxGrowthKb = 64 * 1024 + + proc median3(a, b, c: int64): int64 = + a + b + c - min(a, min(b, c)) - max(a, max(b, c)) + + proc newRssPayload(size: int): string = + result = newString(size) + var index = 0 + while index < size: + result[index] = 'x' + inc index, 4096 + result[^1] = 'x' + + proc drainChannel(iterations, payloadSize: int, useTryRecv: bool): bool = + let channel = newRChan[string](1) + var value: string + for _ in 0 ..< iterations: + channel.send(isolate(newRssPayload(payloadSize))) + if useTryRecv: + if not channel.tryRecv(value): + return false + else: + channel.recv(value) + true + + proc rchanRssGrowth(useTryRecv: bool): int64 = + var samples: array[RssSampleCount, int64] + for index in 0 ..< RssSampleCount: + doAssert drainChannel(RssIterationsPerSample, RssPayloadSize, useTryRecv) + samples[index] = rssKb() + let early = median3(samples[0], samples[1], samples[2]) + let late = median3(samples[3], samples[4], samples[5]) + max(0'i64, late - early) + + suite "RChan managed payload RSS": + test "tryRecv does not retain managed payloads": + check rchanRssGrowth(useTryRecv = true) <= RssMaxGrowthKb + + test "recv does not retain managed payloads": + check rchanRssGrowth(useTryRecv = false) <= RssMaxGrowthKb From 82c840fdbc5ed6023143ba19adbc0662a912fdf0 Mon Sep 17 00:00:00 2001 From: Jaremy Creechley Date: Mon, 21 Sep 2026 19:13:45 +0300 Subject: [PATCH 2/2] use Sigils RChan for image messages --- CHANGES.md | 14 +- figdraw.nimble | 1 + src/figdraw/common/imgutils.nim | 2 +- src/figdraw/common/rchannels.nim | 474 +------------------------------ src/figdraw/commons.nim | 2 +- tests/config.nims | 1 - tests/timage_loading.nim | 48 ++++ tests/trchannels.nim | 356 ----------------------- 8 files changed, 66 insertions(+), 832 deletions(-) delete mode 100644 tests/trchannels.nim diff --git a/CHANGES.md b/CHANGES.md index 279697b8..e5155b52 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -2,10 +2,16 @@ ## 0.41.0 -- Fix managed-value ownership in `RChan` by storing queued values in typed - storage, including overwritten ring-buffer entries and unread messages. -- Release image-message subscription channels and run payload destructors - outside the channel lock, including the `tryTake` commit-race rollback. +- Move the ownership-safe fixed-ring RChan implementation and its regression + suite to Sigils 0.31.0; keep a compatibility import for existing FigDraw + callers and exercise image-message replacement through the shared channel. + +- Store `RChan` messages in a fixed typed ring, transferring isolated values + with `swap` under the lock. Remove per-message allocations and send rollback; + preserve `tryTake` sources on failure and receive isolated values directly. + Initialize vacant slots without constructing payload field defaults. +- Release image-message subscription channels and all ring storage. Run payload + hooks outside channel locks and reject reentrant operations during teardown. - Add Linux/macOS RSS regressions covering both `tryRecv` and blocking `recv`. ## 0.40.2 diff --git a/figdraw.nimble b/figdraw.nimble index 8469f501..6ce4a8d7 100644 --- a/figdraw.nimble +++ b/figdraw.nimble @@ -6,6 +6,7 @@ srcDir = "src" # Dependencies requires "nim >= 2.2" +requires "https://github.com/elcritch/sigils#feat/rchannels" requires "pixie >= 5.0.1" requires "chroma >= 0.2.7" requires "bumpy" diff --git a/src/figdraw/common/imgutils.nim b/src/figdraw/common/imgutils.nim index def78b27..5bc3ff6a 100644 --- a/src/figdraw/common/imgutils.nim +++ b/src/figdraw/common/imgutils.nim @@ -4,7 +4,7 @@ import std/[isolation, locks, times] import pkg/pixie import chronicles -import ./rchannels +import sigils/rchannels import ./formatflippy import ./fonttypes import ./shared diff --git a/src/figdraw/common/rchannels.nim b/src/figdraw/common/rchannels.nim index d08632ee..521e1ac9 100644 --- a/src/figdraw/common/rchannels.nim +++ b/src/figdraw/common/rchannels.nim @@ -1,472 +1,8 @@ -# -# -# Nim's Runtime Library -# (c) Copyright 2021 Andreas Prell, Mamy André-Ratsimbazafy & Nim Contributors -# -# See the file "copying.txt", included in this -# distribution, for details about the copyright. -# -# This Channel implementation is a shared memory, fixed-size, concurrent queue using -# a circular buffer for data. Based on channels implementation[1]_ by -# Mamy André-Ratsimbazafy (@mratsim), which is a C to Nim translation of the -# original[2]_ by Andreas Prell (@aprell) -# -# .. [1] https://github.com/mratsim/weave/blob/5696d94e6358711e840f8c0b7c684fcc5cbd4472/unused/channels/channels_legacy.nim -# .. [2] https://github.com/aprell/tasking-2.0/blob/master/src/channel_shm/channel.c - -## This module works only with one of `--mm:arc` / `--mm:atomicArc` / `--mm:orc` -## compilation flags. -## -## .. warning:: This module is experimental and its interface may change. -## -## This module implements multi-producer multi-consumer channels - a concurrency -## primitive with a high-level interface intended for communication and -## synchronization between threads. It allows sending and receiving typed, isolated -## data, enabling safe and efficient concurrency. -## -## The `RChan` type represents a generic fixed-size channel object that internally manages -## the underlying resources and synchronization. It has to be initialized using -## the `newRChan` proc. Sending and receiving operations are provided by the -## blocking `send` and `recv` procs, and non-blocking `trySend` and `tryRecv` -## procs. For ring buffer behavior, use the `push` proc rather than `send`. -## Send operations add messages to the channel, receiving operations remove them, -## while `push` adds a message or overwrites the oldest message if the channel is full. +## Compatibility import for the RChan implementation now maintained by Sigils. ## -## -## See also: -## * [std/isolation](https://nim-lang.org/docs/isolation.html) -## -## The following is a simple example of two different ways to use channels: -## blocking and non-blocking. - -runnableExamples("--threads:on --gc:orc"): - import std/os - - # In this example a channel is declared at module scope. - # Channels are generic, and they include support for passing objects between - # threads. - # Note that isolated data passed through channels is moved around. - var RChan = newRChan[string]() - - block example_blocking: - # This proc will be run in another thread. - proc basicWorker() = - RChan.send("Hello World!") - - # Launch the worker. - var worker: Thread[void] - createThread(worker, basicWorker) - - # Block until the message arrives, then print it out. - var dest = "" - dest = RChan.recv() - assert dest == "Hello World!" - - # Wait for the thread to exit before moving on to the next example. - worker.joinThread() - - block example_non_blocking: - # This is another proc to run in a background thread. This proc takes a while - # to send the message since it first sleeps for some time. - proc slowWorker(delay: Natural) = - # `delay` is a period in milliseconds - sleep(delay) - RChan.send("Another message") - - # Launch the worker with a delay set to 2 seconds (2000 ms). - var worker: Thread[Natural] - createThread(worker, slowWorker, 2000) - - # This time, use a non-blocking approach with tryRecv. - # Since the main thread is not blocked, it could be used to perform other - # useful work while it waits for data to arrive on the channel. - var messages: seq[string] - while true: - var msg = "" - if RChan.tryRecv(msg): - messages.add msg # "Another message" - break - messages.add "Pretend I'm doing useful work..." - # For this example, sleep in order not to flood the sequence with too many - # "pretend" messages. - sleep(400) - - # Wait for the second thread to exit before cleaning up the channel. - worker.joinThread() - - # Thread exits right after receiving the message - assert messages[^1] == "Another message" - # At least one non-successful attempt to receive the message had to occur. - assert messages.len >= 2 - - block example_non_blocking_overwrite: - var chanRingBuffer = newRChan[string](elements = 1) - chanRingBuffer.push("Hello") - chanRingBuffer.push("World") - var msg = "" - assert chanRingBuffer.tryRecv(msg) - assert msg == "World" - -when not (defined(gcArc) or defined(gcOrc) or defined(gcAtomicArc) or defined(nimdoc)): - {. - error: - "This module requires one of --mm:arc / --mm:atomicArc / --mm:orc compilation flags" - .} - -import std/[atomics, deques, isolation, locks] - -# Channel -# ------------------------------------------------------------------------------ - -type - RChanItemData[T] = object - value: T - - RChanItem[T] = ptr RChanItemData[T] - - RChanData[T] = object - lock: Lock - spaceAvailableCV, dataAvailableCV: Cond - items: Deque[RChanItem[T]] - capacity: int - pendingSends: int - atomicCounter: Atomic[int] - - RChan*[T] = object ## Typed channel - d: ptr RChanData[T] - -when defined(figdrawRChanTests): - var rchanLiveItems*: Atomic[int] - -# The old implementation stored values in an untyped byte buffer and used -# copyMem to move them. That is not valid for managed values: the compiler never -# sees the slot's references, so overwriting or receiving a value cannot destroy -# the old owner. Keep the queue's entries as pointers to individually owned, -# typed values. Pointer operations are safe under the lock; each pointed-to -# value is constructed and destroyed outside the lock so user destructors can -# call back into the channel. - -proc allocItem[T](value: var Isolated[T]): RChanItem[T] = - result = cast[RChanItem[T]](allocShared0(sizeof(RChanItemData[T]))) - when defined(figdrawRChanTests): - discard rchanLiveItems.fetchAdd(1, moRelaxed) - try: - result[].value = extract(value) - except: - try: - `=destroy`(result[].value) - finally: - deallocShared(result) - when defined(figdrawRChanTests): - discard rchanLiveItems.fetchSub(1, moRelaxed) - raise - -proc freeItem[T](item: RChanItem[T]) = - if item.isNil: - return - try: - `=destroy`(item[].value) - finally: - deallocShared(item) - when defined(figdrawRChanTests): - discard rchanLiveItems.fetchSub(1, moRelaxed) - -proc allocChannel[T](n: Positive): ptr RChanData[T] = - result = cast[ptr RChanData[T]](allocShared0(sizeof(RChanData[T]))) - result[].items = initDeque[RChanItem[T]]() - result[].capacity = n - result[].atomicCounter.store(1, moRelaxed) - initLock(result[].lock) - initCond(result[].spaceAvailableCV) - initCond(result[].dataAvailableCV) - -proc freeChannel[T](channel: ptr RChanData[T]) = - if channel.isNil: - return - - # RChanData lives in shared raw storage, so its managed queue field and every - # pointed-to managed value need explicit destruction before the allocation is - # released. No channel lock is held while payload destructors run. - # Detach the queue first: a payload destructor is allowed to re-enter the - # channel, and must not be able to receive an item that teardown is already - # destroying. - var items = channel[].items - channel[].items = initDeque[RChanItem[T]]() - for item in items.items: - freeItem(item) - deinitCond(channel[].spaceAvailableCV) - deinitCond(channel[].dataAvailableCV) - deinitLock(channel[].lock) - deallocShared(channel) - -proc channelSend[T]( - channel: ptr RChanData[T], - value: var Isolated[T], - blocking: static bool, - overwrite: static bool, -): bool = - assert not channel.isNil - - when overwrite: - let overwriteItem = allocItem(value) - var dropped: RChanItem[T] - var overwriteEnqueued = false - try: - acquire(channel[].lock) - try: - if channel[].items.len == channel[].capacity: - # Move the old owner out while holding the lock, but free its value - # after the lock is released. User destructors may call back into the - # channel. - dropped = channel[].items.popFirst() - channel[].items.addLast(overwriteItem) - overwriteEnqueued = true - signal(channel[].dataAvailableCV) - finally: - release(channel[].lock) - except: - if not overwriteEnqueued: - freeItem(overwriteItem) - freeItem(dropped) - raise - freeItem(dropped) - return true - else: - acquire(channel[].lock) - if blocking: - while channel[].items.len + channel[].pendingSends >= channel[].capacity: - wait(channel[].spaceAvailableCV, channel[].lock) - else: - if channel[].items.len + channel[].pendingSends >= channel[].capacity: - release(channel[].lock) - return false - inc channel[].pendingSends - - release(channel[].lock) - var item: RChanItem[T] - try: - item = allocItem(value) - except: - acquire(channel[].lock) - dec channel[].pendingSends - signal(channel[].spaceAvailableCV) - release(channel[].lock) - raise - - var commitRejected = false - var sendEnqueued = false - try: - acquire(channel[].lock) - try: - dec channel[].pendingSends - if not blocking: - if channel[].items.len == channel[].capacity: - commitRejected = true - else: - while channel[].items.len == channel[].capacity: - wait(channel[].spaceAvailableCV, channel[].lock) - if not commitRejected: - channel[].items.addLast(item) - sendEnqueued = true - signal(channel[].dataAvailableCV) - finally: - release(channel[].lock) - except: - if not sendEnqueued: - freeItem(item) - raise - - if commitRejected: - try: - value = unsafeIsolate(move item[].value) - finally: - freeItem(item) - return false - result = true - -proc channelReceive[T]( - channel: ptr RChanData[T], value: var T, blocking: static bool -): bool = - assert not channel.isNil - - var received: RChanItem[T] - acquire(channel[].lock) - when blocking: - while channel[].items.len == 0: - wait(channel[].dataAvailableCV, channel[].lock) - else: - if channel[].items.len == 0: - release(channel[].lock) - return false - - # Move the owner out while holding the lock, but assign it to the caller's - # destination after unlocking. That assignment destroys any value already - # owned by the destination; T's destructor could call back into the channel. - received = channel[].items.popFirst() - signal(channel[].spaceAvailableCV) - release(channel[].lock) - try: - value = move received[].value - finally: - freeItem(received) - result = true - -# Public API -# ------------------------------------------------------------------------------ - -template frees(c) = - if c.d != nil: - # this `fetchSub` returns current val then subs - # fetchSub returns the count before decrementing, so one means this is the - # final owner. - if c.d.atomicCounter.fetchSub(1, moAcquireRelease) == 1: - freeChannel(c.d) - -when defined(nimAllowNonVarDestructor): - proc `=destroy`*[T](c: RChan[T]) = - frees(c) - -else: - proc `=destroy`*[T](c: var RChan[T]) = - frees(c) - -proc `=wasMoved`*[T](x: var RChan[T]) = - x.d = nil - -proc `=dup`*[T](src: RChan[T]): RChan[T] = - if src.d != nil: - discard fetchAdd(src.d.atomicCounter, 1, moRelaxed) - result.d = src.d - -proc `=copy`*[T](dest: var RChan[T], src: RChan[T]) = - ## Shares `Channel` by reference counting. - if src.d != nil: - discard fetchAdd(src.d.atomicCounter, 1, moRelaxed) - `=destroy`(dest) - dest.d = src.d - -proc trySend*[T](c: RChan[T], src: sink Isolated[T]): bool {.inline.} = - ## Tries to send the message `src` to the channel `c`. - ## - ## The memory of `src` will be moved if possible. - ## Doesn't block waiting for space in the channel to become available. - ## Instead returns after an attempt to send a message was made. - ## - ## .. warning:: In high-concurrency situations, consider using an exponential - ## backoff strategy to reduce contention and improve the success rate of - ## operations. - ## - ## Returns `false` if the message was not sent because the number of pending - ## messages in the channel exceeded its capacity. - result = channelSend(c.d, src, false, false) - if result: - wasMoved(src) - -template trySend*[T](c: RChan[T], src: T): bool = - ## Helper template for `trySend <#trySend,RChan[T],sinkIsolated[T]>`_. - ## - ## .. warning:: For repeated sends of the same value, consider using the - ## `tryTake <#tryTake,RChan[T],Isolated[T]>`_ proc with a pre-isolated - ## value to avoid unnecessary copying. - mixin isolate - trySend(c, isolate(src)) - -proc tryTake*[T](c: RChan[T], src: var Isolated[T]): bool {.inline.} = - ## Tries to send the message `src` to the channel `c`. - ## - ## The memory of `src` is moved directly. Be careful not to reuse `src` afterwards. - ## This proc is suitable when `src` cannot be copied. - ## - ## Doesn't block waiting for space in the channel to become available. - ## Instead returns after an attempt to send a message was made. - ## - ## .. warning:: In high-concurrency situations, consider using an exponential - ## backoff strategy to reduce contention and improve the success rate of - ## operations. - ## - ## Returns `false` if the message was not sent because the number of pending - ## messages in the channel exceeded its capacity. - result = channelSend(c.d, src, false, false) - if result: - wasMoved(src) - -proc tryRecv*[T](c: RChan[T], dst: var T): bool {.inline.} = - ## Tries to receive a message from the channel `c` and fill `dst` with its value. - ## - ## Doesn't block waiting for messages in the channel to become available. - ## Instead returns after an attempt to receive a message was made. - ## - ## .. warning:: In high-concurrency situations, consider using an exponential - ## backoff strategy to reduce contention and improve the success rate of - ## operations. - ## - ## Returns `false` and does not change `dst` if no message was received. - channelReceive(c.d, dst, false) - -proc send*[T](c: RChan[T], src: sink Isolated[T]) {.inline.} = - ## Sends the message `src` to the channel `c`. - ## This blocks the sending thread until `src` was successfully sent. - ## - ## The memory of `src` is moved, not copied. - ## - ## If the channel is already full with messages this will block the thread until - ## messages from the channel are removed. - when defined(gcOrc) and defined(nimSafeOrcSend): - GC_runOrc() - discard channelSend(c.d, src, true, false) - wasMoved(src) - -template send*[T](c: RChan[T], src: T) = - ## Helper template for `send`. - mixin isolate - send(c, isolate(src)) - -proc push*[T](c: RChan[T], src: sink Isolated[T]) {.inline.} = - ## Pushes the message `src` to the channel `c`. - ## This is a non-blocking operation that overwrites the oldest message if the channel is full. - ## - ## The memory of `src` is moved, not copied. - when defined(gcOrc) and defined(nimSafeOrcSend): - GC_runOrc() - discard channelSend(c.d, src, false, true) - wasMoved(src) - -template push*[T](c: RChan[T], src: T) = - ## Helper template for `push`. - mixin isolate - push(c, isolate(src)) - -proc recv*[T](c: RChan[T], dst: var T) {.inline.} = - ## Receives a message from the channel `c` and fill `dst` with its value. - ## - ## This blocks the receiving thread until a message was successfully received. - ## - ## If the channel does not contain any messages this will block the thread until - ## a message get sent to the channel. - discard channelReceive(c.d, dst, true) - -proc recv*[T](c: RChan[T]): T {.inline.} = - ## Receives a message from the channel. - ## A version of `recv`_ that returns the message. - discard channelReceive(c.d, result, true) - -proc recvIso*[T](c: RChan[T]): Isolated[T] {.inline.} = - ## Receives a message from the channel. - ## A version of `recv`_ that returns the message and isolates it. - var value: T - discard channelReceive(c.d, value, true) - result = unsafeIsolate(move value) +## New code should import `sigils/rchannels` directly. This module remains so +## existing FigDraw imports continue to resolve while callers migrate. -proc peek*[T](c: RChan[T]): int {.inline.} = - ## Returns an estimation of the current number of messages held by the channel. - acquire(c.d[].lock) - result = c.d[].items.len - release(c.d[].lock) +import sigils/rchannels -proc newRChan*[T](elements: Positive = 30): RChan[T] = - ## An initialization procedure, necessary for acquiring resources and - ## initializing internal state of the channel. - ## - ## `elements` is the capacity of the channel and thus how many messages it can hold - ## before it refuses to accept any further messages. - result = RChan[T](d: allocChannel[T](elements)) +export rchannels diff --git a/src/figdraw/commons.nim b/src/figdraw/commons.nim index 1307a657..8a15192a 100644 --- a/src/figdraw/commons.nim +++ b/src/figdraw/commons.nim @@ -1,6 +1,6 @@ import common/shared import common/uimaths -import common/rchannels +import sigils/rchannels import common/fontutils import common/imgutils import extras/systemfonts diff --git a/tests/config.nims b/tests/config.nims index 3522bf07..5fa5d1c5 100644 --- a/tests/config.nims +++ b/tests/config.nims @@ -1,3 +1,2 @@ --path:"../src" --define:"figdraw.vulkanReadback" ---define:"figdrawRChanTests" diff --git a/tests/timage_loading.nim b/tests/timage_loading.nim index c6d4faf7..55dfa4f0 100644 --- a/tests/timage_loading.nim +++ b/tests/timage_loading.nim @@ -124,6 +124,13 @@ proc retainImageOnThread(id: ImageId) {.thread.} = var owned = imageRef(id) discard owned.id +when not defined(useMalloc): + proc createUnreadSubscription() = + let subscription = newImageMessageSubscription() + doAssert not subscription.isNil + # Leave the replayed image in the inbox so subscription teardown must + # release both the ring storage and its managed payload. + suite "image loading": test "load png via figDataDir fallback": setFigDataDir(getCurrentDir() / "data") @@ -266,6 +273,33 @@ suite "image loading": ctx.drainImages() late.drainImages() + test "image replacement exercises the bounded RChan ring": + clearImageCache() + let + id = imgId("tests/timage_loading/rchan-ring") + subscription = newImageMessageSubscription() + frameCount = 4300 + for frame in 0 ..< frameCount: + var image = newImage(1, 1) + image[0, 0] = rgba(uint8(frame mod 251), 20, 30, 255) + if frame == 0: + loadImage(id, image) + else: + replaceImage(id, image) + + var + message: ImageMsg + received = 0 + lastRed: uint8 + while subscription.tryRecvImageMsg(message): + if message.id == id and message.kind in {ImkPutPixie, ImkReplacePixie}: + lastRed = message.pimg[0, 0].r + inc received + + check received > 0 + check lastRed == uint8((frameCount - 1) mod 251) + clearImage(id) + test "replaceImage replays the newest frame after an atlas rebuild": let id = imgId("tests/timage_loading/replace/rebuild") @@ -579,6 +613,20 @@ suite "image loading": clearImage(id) + when not defined(useMalloc): + test "subscription teardown releases its inbox and unread image replay": + clearImageCache() + let id = imgId("tests/timage_loading/subscription-cleanup") + loadImage(id, newImage(64, 64)) + # Warm up the subscriber table before comparing live allocations. + createUnreadSubscription() + let before = getOccupiedMem() + for _ in 0 ..< 100: + createUnreadSubscription() + let growth = getOccupiedMem() - before + check growth == 0 + clearImage(id) + test "a full renderer subscription does not block resource publishers": clearImageCache() let subscription = newImageMessageSubscription() diff --git a/tests/trchannels.nim b/tests/trchannels.nim deleted file mode 100644 index 7e505791..00000000 --- a/tests/trchannels.nim +++ /dev/null @@ -1,356 +0,0 @@ -import std/[atomics, isolation, os, unittest] -import figdraw/common/rchannels - -var destroyed: Atomic[int] - -type Payload = object - id: int - -proc `=destroy`(payload: Payload) = - if payload.id != 0: - discard destroyed.fetchAdd(1) - -type ReentrantPayload = object - id: int - -var - reentrantChannel: RChan[ReentrantPayload] - reentrantAlias: RChan[ReentrantPayload] - reentrantReady: bool - reentrantReceive: bool - reentrantCallbackDone: bool - -proc `=destroy`(payload: ReentrantPayload) = - if reentrantReady and not reentrantCallbackDone: - reentrantCallbackDone = true - if reentrantReceive: - var ignored: ReentrantPayload - discard reentrantAlias.tryRecv(ignored) - else: - discard reentrantChannel.peek() - -type BlockingPayload = object - value: int - -type RaisingPayload = object - id: int - -var raiseOnSink: bool - -proc `=sink`(dest: var RaisingPayload, src: RaisingPayload) = - if raiseOnSink: - raise newException(ValueError, "test sink failure") - dest.id = src.id - -var - blockingSinkStarted: Atomic[bool] - allowBlockingSink: Atomic[bool] - blockingSinkRaise: bool - raceInitialDone: Atomic[bool] - allowRaceRetry: Atomic[bool] - raceResult: Atomic[bool] - raceRetryResult: Atomic[bool] - raceSource: Atomic[int] - -proc `=sink`(dest: var BlockingPayload, src: BlockingPayload) = - if blockingSinkRaise: - raise newException(ValueError, "rollback sink failure") - dest.value = src.value - if not blockingSinkStarted.load(moAcquire): - blockingSinkStarted.store(true, moRelease) - while not allowBlockingSink.load(moAcquire): - sleep(1) - -var raceChannel: RChan[BlockingPayload] - -proc tryTakeDuringPush() {.thread.} = - {.cast(gcsafe).}: - var source = isolate(BlockingPayload(value: 42)) - let succeeded = raceChannel.tryTake(source) - raceResult.store(succeeded, moRelease) - if not succeeded: - var restored = extract(source) - raceSource.store(restored.value, moRelease) - source = isolate(move restored) - raceInitialDone.store(true, moRelease) - if not succeeded: - while not allowRaceRetry.load(moAcquire): - sleep(1) - raceRetryResult.store(raceChannel.tryTake(source), moRelease) - -proc tryTakeWithRaisingRollback() {.thread.} = - {.cast(gcsafe).}: - var source = isolate(BlockingPayload(value: 42)) - try: - discard raceChannel.tryTake(source) - except ValueError: - raceRetryResult.store(true, moRelease) - -const RaceWaitLimit = 5000 - -proc waitForFlag(flag: var Atomic[bool]): bool = - for _ in 0 ..< RaceWaitLimit: - if flag.load(moAcquire): - return true - sleep(1) - -suite "RChan managed payload ownership": - test "tryRecv destroys received payloads exactly once": - const iterations = 100 - destroyed.store(0) - block: - let channel = newRChan[Payload](1) - var value: Payload - for id in 1 .. iterations: - channel.send(isolate(Payload(id: id))) - check channel.tryRecv(value) - check value.id == id - check destroyed.load() == iterations - - test "recv destroys received payloads exactly once": - const iterations = 100 - destroyed.store(0) - block: - let channel = newRChan[Payload](1) - for id in 1 .. iterations: - channel.send(isolate(Payload(id: id))) - let value = channel.recv() - check value.id == id - check destroyed.load() == iterations - - test "recv into a destination destroys the previous value": - const iterations = 100 - destroyed.store(0) - block: - let channel = newRChan[Payload](1) - var value: Payload - for id in 1 .. iterations: - channel.send(isolate(Payload(id: id))) - channel.recv(value) - check value.id == id - check destroyed.load() == iterations - - test "push destroys overwritten payloads": - destroyed.store(0) - block: - let channel = newRChan[Payload](1) - channel.push(isolate(Payload(id: 1))) - channel.push(isolate(Payload(id: 2))) - check channel.recv().id == 2 - check destroyed.load() == 2 - - test "the final shared owner destroys unread payloads": - destroyed.store(0) - block: - var retained: RChan[Payload] - block: - let channel = newRChan[Payload](1) - retained = channel - channel.send(isolate(Payload(id: 1))) - check destroyed.load() == 0 - check destroyed.load() == 1 - - test "payload destruction can re-enter the channel": - reentrantReady = false - reentrantReceive = false - reentrantCallbackDone = false - block: - let channel = newRChan[ReentrantPayload](1) - reentrantChannel = channel - reentrantReady = true - channel.send(isolate(ReentrantPayload(id: 1))) - var value: ReentrantPayload - check channel.tryRecv(value) - reentrantReady = false - reset(reentrantChannel) - - test "final-owner cleanup detaches queued items before destruction": - reentrantReady = true - reentrantReceive = true - reentrantCallbackDone = false - block: - let channel = newRChan[ReentrantPayload](1) - copyMem(addr reentrantAlias, unsafeAddr channel, sizeof(reentrantAlias)) - channel.send(isolate(ReentrantPayload(id: 1))) - wasMoved(reentrantAlias) - reentrantReady = false - reentrantReceive = false - - test "tryTake preserves its source when push wins the commit race": - blockingSinkStarted.store(false) - allowBlockingSink.store(false) - raceInitialDone.store(false) - allowRaceRetry.store(false) - raceResult.store(true) - raceRetryResult.store(false) - raceSource.store(0) - raceChannel = newRChan[BlockingPayload](1) - - var worker: Thread[void] - createThread(worker, tryTakeDuringPush) - let sinkStarted = waitForFlag(blockingSinkStarted) - if sinkStarted: - raceChannel.push(BlockingPayload(value: 99)) - allowBlockingSink.store(true, moRelease) - let initialDone = waitForFlag(raceInitialDone) - if initialDone and not raceResult.load(moAcquire): - var queued: BlockingPayload - check raceChannel.tryRecv(queued) - check queued.value == 99 - allowRaceRetry.store(true, moRelease) - joinThread(worker) - - check sinkStarted - check initialDone - if sinkStarted and initialDone: - check not raceResult.load(moAcquire) - check raceSource.load(moAcquire) == 42 - check raceRetryResult.load(moAcquire) - var retried: BlockingPayload - check raceChannel.tryRecv(retried) - check retried.value == 42 - var leftover: BlockingPayload - discard raceChannel.tryRecv(leftover) - raceChannel = RChan[BlockingPayload]() - - test "raising ownership hooks do not strand channel capacity": - let channel = newRChan[RaisingPayload](1) - var source = isolate(RaisingPayload(id: 1)) - raiseOnSink = true - expect ValueError: - discard channel.tryTake(source) - raiseOnSink = false - - check channel.peek() == 0 - channel.send(isolate(RaisingPayload(id: 2))) - check channel.recv().id == 2 - - channel.send(isolate(RaisingPayload(id: 3))) - var destination = RaisingPayload(id: 0) - raiseOnSink = true - expect ValueError: - channel.recv(destination) - raiseOnSink = false - - check channel.peek() == 0 - channel.send(isolate(RaisingPayload(id: 4))) - check channel.recv().id == 4 - - test "raising rollback hooks still free the rejected item": - let liveItemsBefore = rchanLiveItems.load(moAcquire) - blockingSinkStarted.store(false) - allowBlockingSink.store(false) - blockingSinkRaise = false - raceRetryResult.store(false) - raceChannel = newRChan[BlockingPayload](1) - - var worker: Thread[void] - createThread(worker, tryTakeWithRaisingRollback) - let sinkStarted = waitForFlag(blockingSinkStarted) - if sinkStarted: - raceChannel.push(BlockingPayload(value: 99)) - blockingSinkRaise = true - allowBlockingSink.store(true, moRelease) - joinThread(worker) - blockingSinkRaise = false - - check sinkStarted - check raceRetryResult.load(moAcquire) - var queued: BlockingPayload - check raceChannel.tryRecv(queued) - check queued.value == 99 - raceChannel = RChan[BlockingPayload]() - check rchanLiveItems.load(moAcquire) == liveItemsBefore - -when defined(linux): - import std/strutils - - proc rssKb(): int64 = - for line in lines("/proc/self/status"): - if line.startsWith("VmRSS:"): - let fields = line.splitWhitespace() - if fields.len >= 2: - return parseInt(fields[1]).int64 - raise newException(IOError, "unable to read VmRSS from /proc/self/status") - -elif defined(macosx): - type - MachPort = uint32 - MachMsgTypeNumber = uint32 - TimeValue = object - seconds: int32 - microseconds: int32 - - MachTaskBasicInfo = object - virtualSize: uint64 - residentSize: uint64 - residentSizeMax: uint64 - userTime: TimeValue - systemTime: TimeValue - policy: int32 - suspendCount: int32 - - const machTaskBasicInfoFlavor = 20 - - var machTaskSelf {.importc: "mach_task_self_", header: "".}: - MachPort - - proc taskInfo( - task: MachPort, - flavor: cint, - taskInfoOut: ptr MachTaskBasicInfo, - taskInfoOutCount: ptr MachMsgTypeNumber, - ): cint {.importc: "task_info", header: "".} - - proc rssKb(): int64 = - var info: MachTaskBasicInfo - var count = MachMsgTypeNumber(sizeof(MachTaskBasicInfo) div sizeof(cuint)) - let status = taskInfo(machTaskSelf, machTaskBasicInfoFlavor, addr info, addr count) - doAssert status == 0, "task_info failed with kern_return_t " & $status - int64(info.residentSize div 1024) - -when defined(linux) or defined(macosx): - const - RssSampleCount = 6 - RssIterationsPerSample = 64 - RssPayloadSize = 512 * 1024 - RssMaxGrowthKb = 64 * 1024 - - proc median3(a, b, c: int64): int64 = - a + b + c - min(a, min(b, c)) - max(a, max(b, c)) - - proc newRssPayload(size: int): string = - result = newString(size) - var index = 0 - while index < size: - result[index] = 'x' - inc index, 4096 - result[^1] = 'x' - - proc drainChannel(iterations, payloadSize: int, useTryRecv: bool): bool = - let channel = newRChan[string](1) - var value: string - for _ in 0 ..< iterations: - channel.send(isolate(newRssPayload(payloadSize))) - if useTryRecv: - if not channel.tryRecv(value): - return false - else: - channel.recv(value) - true - - proc rchanRssGrowth(useTryRecv: bool): int64 = - var samples: array[RssSampleCount, int64] - for index in 0 ..< RssSampleCount: - doAssert drainChannel(RssIterationsPerSample, RssPayloadSize, useTryRecv) - samples[index] = rssKb() - let early = median3(samples[0], samples[1], samples[2]) - let late = median3(samples[3], samples[4], samples[5]) - max(0'i64, late - early) - - suite "RChan managed payload RSS": - test "tryRecv does not retain managed payloads": - check rchanRssGrowth(useTryRecv = true) <= RssMaxGrowthKb - - test "recv does not retain managed payloads": - check rchanRssGrowth(useTryRecv = false) <= RssMaxGrowthKb