20#include <winpr/config.h>
23#include <winpr/wlog.h>
25#include <winpr/collections.h>
29#define XTAG WINPR_TAG("utils.streampool")
31#if !defined(STREAMPOOL_SIZE_LIMIT)
32#error "CMake must define STREAMPOOL_SIZE_LIMIT=<unsigned>"
34static const size_t POOL_COMMON_LIMIT = STREAMPOOL_SIZE_LIMIT;
36struct s_StreamPoolEntry
38#if defined(WITH_DEBUG_STREAMPOOL)
49 struct s_StreamPoolEntry* aArray;
53 struct s_StreamPoolEntry* uArray;
61static void discard_entry(
struct s_StreamPoolEntry* entry, BOOL discardStream)
66#if defined(WITH_DEBUG_STREAMPOOL)
67 free((
void*)entry->msg);
70 if (discardStream && entry->s)
71 Stream_Free(entry->s, entry->s->isAllocatedStream);
73 const struct s_StreamPoolEntry empty = WINPR_C_ARRAY_INIT;
77static struct s_StreamPoolEntry add_entry(
wStream* s)
79 struct s_StreamPoolEntry entry = WINPR_C_ARRAY_INIT;
81#if defined(WITH_DEBUG_STREAMPOOL)
82 void* stack = winpr_backtrace(20);
84 entry.msg = winpr_backtrace_symbols(stack, &entry.lines);
85 winpr_backtrace_free(stack);
96static inline void StreamPool_Lock(wStreamPool* pool)
99 if (pool->synchronized)
100 EnterCriticalSection(&pool->lock);
107static inline void StreamPool_Unlock(wStreamPool* pool)
110 if (pool->synchronized)
111 LeaveCriticalSection(&pool->lock);
114static BOOL StreamPool_ShrinkToCommonLimit(
wStream* s)
116 if (Stream_Capacity(s) <= POOL_COMMON_LIMIT)
119 return Stream_ResizeToCapacity(s, POOL_COMMON_LIMIT);
122static BOOL StreamPool_EnsureCapacity(wStreamPool* pool,
size_t count, BOOL usedOrAvailable)
126 size_t* cap = (usedOrAvailable) ? &pool->uCapacity : &pool->aCapacity;
127 size_t* size = (usedOrAvailable) ? &pool->uSize : &pool->aSize;
128 struct s_StreamPoolEntry** array = (usedOrAvailable) ? &pool->uArray : &pool->aArray;
132 new_cap = *size + count;
133 else if (*size + count > *cap)
134 new_cap = (*size + count + 2) / 2 * 3;
135 else if ((*size + count) < *cap / 3)
140 struct s_StreamPoolEntry* new_arr =
nullptr;
142 if (*cap < *size + count)
146 (
struct s_StreamPoolEntry*)realloc(*array,
sizeof(
struct s_StreamPoolEntry) * new_cap);
159static void StreamPool_ShiftUsed(wStreamPool* pool,
size_t index)
163 const size_t pcount = 1;
164 const size_t off = index + pcount;
165 if (pool->uSize >= off)
167 for (
size_t x = 0; x < pcount; x++)
169 struct s_StreamPoolEntry* cur = &pool->uArray[index + x];
170 discard_entry(cur, FALSE);
172 MoveMemory(&pool->uArray[index], &pool->uArray[index + pcount],
173 (pool->uSize - index - pcount) *
sizeof(
struct s_StreamPoolEntry));
174 pool->uSize -= pcount;
182static void StreamPool_AddUsed(wStreamPool* pool,
wStream* s)
184 StreamPool_EnsureCapacity(pool, 1, TRUE);
185 pool->uArray[pool->uSize] = add_entry(s);
193static void StreamPool_RemoveUsed(wStreamPool* pool,
wStream* s)
196 for (
size_t index = 0; index < pool->uSize; index++)
198 struct s_StreamPoolEntry* cur = &pool->uArray[index];
201 StreamPool_ShiftUsed(pool, index);
207static void StreamPool_ShiftAvailable(wStreamPool* pool,
size_t index)
211 const size_t pcount = 1;
212 const size_t off = index + pcount;
213 if (pool->aSize >= off)
215 for (
size_t x = 0; x < pcount; x++)
217 struct s_StreamPoolEntry* cur = &pool->aArray[index + x];
218 discard_entry(cur, FALSE);
221 MoveMemory(&pool->aArray[index], &pool->aArray[index + pcount],
222 (pool->aSize - index - pcount) *
sizeof(
struct s_StreamPoolEntry));
223 pool->aSize -= pcount;
231wStream* StreamPool_Take(wStreamPool* pool,
size_t size)
234 size_t foundIndex = 0;
237 StreamPool_Lock(pool);
240 size = pool->defaultSize;
242 for (
size_t index = 0; index < pool->aSize; index++)
244 struct s_StreamPoolEntry* cur = &pool->aArray[index];
247 if (Stream_Capacity(s) >= size)
257 s = Stream_New(
nullptr, size);
263 Stream_ResetPosition(s);
264 if (!Stream_SetLength(s, Stream_Capacity(s)))
266 StreamPool_ShiftAvailable(pool, foundIndex);
273 StreamPool_AddUsed(pool, s);
277 StreamPool_Unlock(pool);
286static void StreamPool_Remove(wStreamPool* pool,
wStream* s)
288 StreamPool_EnsureCapacity(pool, 1, FALSE);
289 Stream_EnsureValidity(s);
290 StreamPool_ShrinkToCommonLimit(s);
291 for (
size_t x = 0; x < pool->aSize; x++)
293 wStream* cs = pool->aArray[x].s;
297 pool->aArray[(pool->aSize)++] = add_entry(s);
298 StreamPool_RemoveUsed(pool, s);
301static void StreamPool_ReleaseOrReturn(wStreamPool* pool,
wStream* s)
303 StreamPool_Lock(pool);
304 StreamPool_Remove(pool, s);
305 StreamPool_Unlock(pool);
308void StreamPool_Return(wStreamPool* pool,
wStream* s)
314 StreamPool_Lock(pool);
315 StreamPool_Remove(pool, s);
316 StreamPool_Unlock(pool);
342 StreamPool_ReleaseOrReturn(s->pool, s);
344 Stream_Free(s, TRUE);
352wStream* StreamPool_Find(wStreamPool* pool,
const BYTE* ptr)
356 StreamPool_Lock(pool);
358 for (
size_t index = 0; index < pool->uSize; index++)
360 struct s_StreamPoolEntry* cur = &pool->uArray[index];
362 if ((ptr >= Stream_Buffer(cur->s)) &&
363 (ptr < (Stream_Buffer(cur->s) + Stream_Capacity(cur->s))))
370 StreamPool_Unlock(pool);
379void StreamPool_Clear(wStreamPool* pool)
381 StreamPool_Lock(pool);
383 for (
size_t x = 0; x < pool->aSize; x++)
385 struct s_StreamPoolEntry* cur = &pool->aArray[x];
386 discard_entry(cur, TRUE);
392 WLog_Print(pool->log, WLOG_WARN,
393 "Clearing StreamPool, but there are %" PRIuz
" streams currently in use",
395 for (
size_t x = 0; x < pool->uSize; x++)
397 struct s_StreamPoolEntry* cur = &pool->uArray[x];
398 discard_entry(cur, TRUE);
403 StreamPool_Unlock(pool);
406size_t StreamPool_UsedCount(wStreamPool* pool)
408 StreamPool_Lock(pool);
409 size_t usize = pool->uSize;
410 StreamPool_Unlock(pool);
418wStreamPool* StreamPool_New(BOOL
synchronized,
size_t defaultSize)
420 wStreamPool* pool = calloc(1,
sizeof(wStreamPool));
425 pool->log = WLog_Create(XTAG, WLog_GetRoot());
429 pool->synchronized =
synchronized;
430 pool->defaultSize = defaultSize;
432 if (!StreamPool_EnsureCapacity(pool, 32, FALSE))
434 if (!StreamPool_EnsureCapacity(pool, 32, TRUE))
437 if (!InitializeCriticalSectionAndSpinCount(&pool->lock, 4000))
442 WINPR_PRAGMA_DIAG_PUSH
443 WINPR_PRAGMA_DIAG_IGNORED_MISMATCHED_DEALLOC
444 StreamPool_Free(pool);
445 WINPR_PRAGMA_DIAG_POP
449void StreamPool_Free(wStreamPool* pool)
454 StreamPool_Clear(pool);
456 DeleteCriticalSection(&pool->lock);
461 WLog_Discard(pool->log);
465char* StreamPool_GetStatistics(wStreamPool* pool,
char* buffer,
size_t size)
469 if (!buffer || (size < 1))
473 int offset = _snprintf(buffer, size - 1,
474 "aSize =%" PRIuz
", uSize =%" PRIuz
", aCapacity=%" PRIuz
475 ", uCapacity=%" PRIuz,
476 pool->aSize, pool->uSize, pool->aCapacity, pool->uCapacity);
477 if ((offset > 0) && ((
size_t)offset < size))
478 used += (size_t)offset;
480#if defined(WITH_DEBUG_STREAMPOOL)
481 StreamPool_Lock(pool);
483 offset = _snprintf(&buffer[used], size - 1 - used,
"\n-- dump used array take locations --\n");
484 if ((offset > 0) && ((
size_t)offset < size - used))
485 used += (size_t)offset;
486 for (
size_t x = 0; x < pool->uSize; x++)
488 const struct s_StreamPoolEntry* cur = &pool->uArray[x];
489 WINPR_ASSERT(cur->msg || (cur->lines == 0));
491 for (
size_t y = 0; y < cur->lines; y++)
493 offset = _snprintf(&buffer[used], size - 1 - used,
"[%" PRIuz
" | %" PRIuz
"]: %s\n", x,
495 if ((offset > 0) && ((
size_t)offset < size - used))
496 used += (size_t)offset;
500 offset = _snprintf(&buffer[used], size - 1 - used,
"\n-- statistics called from --\n");
501 if ((offset > 0) && ((
size_t)offset < size - used))
502 used += (size_t)offset;
504 struct s_StreamPoolEntry entry = WINPR_C_ARRAY_INIT;
505 void* stack = winpr_backtrace(20);
507 entry.msg = winpr_backtrace_symbols(stack, &entry.lines);
508 winpr_backtrace_free(stack);
510 for (
size_t x = 0; x < entry.lines; x++)
512 const char* msg = entry.msg[x];
513 offset = _snprintf(&buffer[used], size - 1 - used,
"[%" PRIuz
"]: %s\n", x, msg);
514 if ((offset > 0) && ((
size_t)offset < size - used))
515 used += (size_t)offset;
517 free((
void*)entry.msg);
518 StreamPool_Unlock(pool);
524BOOL StreamPool_WaitForReturn(wStreamPool* pool, UINT32 timeoutMS)
528 while (timeoutMS > 0)
530 const size_t used = StreamPool_UsedCount(pool);
533 WLog_Print(pool->log, WLOG_DEBUG,
"%" PRIuz
" streams still in use, sleeping...", used);
535 char buffer[4096] = WINPR_C_ARRAY_INIT;
536 StreamPool_GetStatistics(pool, buffer,
sizeof(buffer));
537 WLog_Print(pool->log, WLOG_TRACE,
"Pool statistics: %s", buffer);
540 if (timeoutMS != INFINITE)
542 diff = timeoutMS > 10 ? 10 : timeoutMS;