FreeRDP
Loading...
Searching...
No Matches
StreamPool.c
1
20#include <winpr/config.h>
21
22#include <winpr/crt.h>
23#include <winpr/wlog.h>
24
25#include <winpr/collections.h>
26
27#include "../stream.h"
28#include "../log.h"
29#define XTAG WINPR_TAG("utils.streampool")
30
31#if !defined(STREAMPOOL_SIZE_LIMIT)
32#error "CMake must define STREAMPOOL_SIZE_LIMIT=<unsigned>"
33#endif
34static const size_t POOL_COMMON_LIMIT = STREAMPOOL_SIZE_LIMIT;
35
36struct s_StreamPoolEntry
37{
38#if defined(WITH_DEBUG_STREAMPOOL)
39 char** msg;
40 size_t lines;
41#endif
42 wStream* s;
43};
44
45struct s_wStreamPool
46{
47 size_t aSize;
48 size_t aCapacity;
49 struct s_StreamPoolEntry* aArray;
50
51 size_t uSize;
52 size_t uCapacity;
53 struct s_StreamPoolEntry* uArray;
54
56 BOOL synchronized;
57 size_t defaultSize;
58 wLog* log;
59};
60
61static void discard_entry(struct s_StreamPoolEntry* entry, BOOL discardStream)
62{
63 if (!entry)
64 return;
65
66#if defined(WITH_DEBUG_STREAMPOOL)
67 free((void*)entry->msg);
68#endif
69
70 if (discardStream && entry->s)
71 Stream_Free(entry->s, entry->s->isAllocatedStream);
72
73 const struct s_StreamPoolEntry empty = WINPR_C_ARRAY_INIT;
74 *entry = empty;
75}
76
77static struct s_StreamPoolEntry add_entry(wStream* s)
78{
79 struct s_StreamPoolEntry entry = WINPR_C_ARRAY_INIT;
80
81#if defined(WITH_DEBUG_STREAMPOOL)
82 void* stack = winpr_backtrace(20);
83 if (stack)
84 entry.msg = winpr_backtrace_symbols(stack, &entry.lines);
85 winpr_backtrace_free(stack);
86#endif
87
88 entry.s = s;
89 return entry;
90}
91
96static inline void StreamPool_Lock(wStreamPool* pool)
97{
98 WINPR_ASSERT(pool);
99 if (pool->synchronized)
100 EnterCriticalSection(&pool->lock);
101}
102
107static inline void StreamPool_Unlock(wStreamPool* pool)
108{
109 WINPR_ASSERT(pool);
110 if (pool->synchronized)
111 LeaveCriticalSection(&pool->lock);
112}
113
114static BOOL StreamPool_ShrinkToCommonLimit(wStream* s)
115{
116 if (Stream_Capacity(s) <= POOL_COMMON_LIMIT)
117 return TRUE;
118
119 return Stream_ResizeToCapacity(s, POOL_COMMON_LIMIT);
120}
121
122static BOOL StreamPool_EnsureCapacity(wStreamPool* pool, size_t count, BOOL usedOrAvailable)
123{
124 WINPR_ASSERT(pool);
125
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;
129
130 size_t new_cap = 0;
131 if (*cap == 0)
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)
136 new_cap = *cap / 2;
137
138 if (new_cap > 0)
139 {
140 struct s_StreamPoolEntry* new_arr = nullptr;
141
142 if (*cap < *size + count)
143 *cap += count;
144
145 new_arr =
146 (struct s_StreamPoolEntry*)realloc(*array, sizeof(struct s_StreamPoolEntry) * new_cap);
147 if (!new_arr)
148 return FALSE;
149 *cap = new_cap;
150 *array = new_arr;
151 }
152 return TRUE;
153}
154
159static void StreamPool_ShiftUsed(wStreamPool* pool, size_t index)
160{
161 WINPR_ASSERT(pool);
162
163 const size_t pcount = 1;
164 const size_t off = index + pcount;
165 if (pool->uSize >= off)
166 {
167 for (size_t x = 0; x < pcount; x++)
168 {
169 struct s_StreamPoolEntry* cur = &pool->uArray[index + x];
170 discard_entry(cur, FALSE);
171 }
172 MoveMemory(&pool->uArray[index], &pool->uArray[index + pcount],
173 (pool->uSize - index - pcount) * sizeof(struct s_StreamPoolEntry));
174 pool->uSize -= pcount;
175 }
176}
177
182static void StreamPool_AddUsed(wStreamPool* pool, wStream* s)
183{
184 StreamPool_EnsureCapacity(pool, 1, TRUE);
185 pool->uArray[pool->uSize] = add_entry(s);
186 pool->uSize++;
187}
188
193static void StreamPool_RemoveUsed(wStreamPool* pool, wStream* s)
194{
195 WINPR_ASSERT(pool);
196 for (size_t index = 0; index < pool->uSize; index++)
197 {
198 struct s_StreamPoolEntry* cur = &pool->uArray[index];
199 if (cur->s == s)
200 {
201 StreamPool_ShiftUsed(pool, index);
202 break;
203 }
204 }
205}
206
207static void StreamPool_ShiftAvailable(wStreamPool* pool, size_t index)
208{
209 WINPR_ASSERT(pool);
210
211 const size_t pcount = 1;
212 const size_t off = index + pcount;
213 if (pool->aSize >= off)
214 {
215 for (size_t x = 0; x < pcount; x++)
216 {
217 struct s_StreamPoolEntry* cur = &pool->aArray[index + x];
218 discard_entry(cur, FALSE);
219 }
220
221 MoveMemory(&pool->aArray[index], &pool->aArray[index + pcount],
222 (pool->aSize - index - pcount) * sizeof(struct s_StreamPoolEntry));
223 pool->aSize -= pcount;
224 }
225}
226
231wStream* StreamPool_Take(wStreamPool* pool, size_t size)
232{
233 BOOL found = FALSE;
234 size_t foundIndex = 0;
235 wStream* s = nullptr;
236
237 StreamPool_Lock(pool);
238
239 if (size == 0)
240 size = pool->defaultSize;
241
242 for (size_t index = 0; index < pool->aSize; index++)
243 {
244 struct s_StreamPoolEntry* cur = &pool->aArray[index];
245 s = cur->s;
246
247 if (Stream_Capacity(s) >= size)
248 {
249 found = TRUE;
250 foundIndex = index;
251 break;
252 }
253 }
254
255 if (!found)
256 {
257 s = Stream_New(nullptr, size);
258 if (!s)
259 goto out_fail;
260 }
261 else if (s)
262 {
263 Stream_ResetPosition(s);
264 if (!Stream_SetLength(s, Stream_Capacity(s)))
265 goto out_fail;
266 StreamPool_ShiftAvailable(pool, foundIndex);
267 }
268
269 if (s)
270 {
271 s->pool = pool;
272 s->count = 1;
273 StreamPool_AddUsed(pool, s);
274 }
275
276out_fail:
277 StreamPool_Unlock(pool);
278
279 return s;
280}
281
286static void StreamPool_Remove(wStreamPool* pool, wStream* s)
287{
288 StreamPool_EnsureCapacity(pool, 1, FALSE);
289 Stream_EnsureValidity(s);
290 StreamPool_ShrinkToCommonLimit(s);
291 for (size_t x = 0; x < pool->aSize; x++)
292 {
293 wStream* cs = pool->aArray[x].s;
294 if (cs == s)
295 return;
296 }
297 pool->aArray[(pool->aSize)++] = add_entry(s);
298 StreamPool_RemoveUsed(pool, s);
299}
300
301static void StreamPool_ReleaseOrReturn(wStreamPool* pool, wStream* s)
302{
303 StreamPool_Lock(pool);
304 StreamPool_Remove(pool, s);
305 StreamPool_Unlock(pool);
306}
307
308void StreamPool_Return(wStreamPool* pool, wStream* s)
309{
310 WINPR_ASSERT(pool);
311 if (!s)
312 return;
313
314 StreamPool_Lock(pool);
315 StreamPool_Remove(pool, s);
316 StreamPool_Unlock(pool);
317}
318
323void Stream_AddRef(wStream* s)
324{
325 WINPR_ASSERT(s);
326 s->count++;
327}
328
333void Stream_Release(wStream* s)
334{
335 WINPR_ASSERT(s);
336
337 if (s->count > 0)
338 s->count--;
339 if (s->count == 0)
340 {
341 if (s->pool)
342 StreamPool_ReleaseOrReturn(s->pool, s);
343 else
344 Stream_Free(s, TRUE);
345 }
346}
347
352wStream* StreamPool_Find(wStreamPool* pool, const BYTE* ptr)
353{
354 wStream* s = nullptr;
355
356 StreamPool_Lock(pool);
357
358 for (size_t index = 0; index < pool->uSize; index++)
359 {
360 struct s_StreamPoolEntry* cur = &pool->uArray[index];
361
362 if ((ptr >= Stream_Buffer(cur->s)) &&
363 (ptr < (Stream_Buffer(cur->s) + Stream_Capacity(cur->s))))
364 {
365 s = cur->s;
366 break;
367 }
368 }
369
370 StreamPool_Unlock(pool);
371
372 return s;
373}
374
379void StreamPool_Clear(wStreamPool* pool)
380{
381 StreamPool_Lock(pool);
382
383 for (size_t x = 0; x < pool->aSize; x++)
384 {
385 struct s_StreamPoolEntry* cur = &pool->aArray[x];
386 discard_entry(cur, TRUE);
387 }
388 pool->aSize = 0;
389
390 if (pool->uSize > 0)
391 {
392 WLog_Print(pool->log, WLOG_WARN,
393 "Clearing StreamPool, but there are %" PRIuz " streams currently in use",
394 pool->uSize);
395 for (size_t x = 0; x < pool->uSize; x++)
396 {
397 struct s_StreamPoolEntry* cur = &pool->uArray[x];
398 discard_entry(cur, TRUE);
399 }
400 pool->uSize = 0;
401 }
402
403 StreamPool_Unlock(pool);
404}
405
406size_t StreamPool_UsedCount(wStreamPool* pool)
407{
408 StreamPool_Lock(pool);
409 size_t usize = pool->uSize;
410 StreamPool_Unlock(pool);
411 return usize;
412}
413
418wStreamPool* StreamPool_New(BOOL synchronized, size_t defaultSize)
419{
420 wStreamPool* pool = calloc(1, sizeof(wStreamPool));
421
422 if (!pool)
423 return nullptr;
424
425 pool->log = WLog_Create(XTAG, WLog_GetRoot());
426 if (!pool->log)
427 goto fail;
428
429 pool->synchronized = synchronized;
430 pool->defaultSize = defaultSize;
431
432 if (!StreamPool_EnsureCapacity(pool, 32, FALSE))
433 goto fail;
434 if (!StreamPool_EnsureCapacity(pool, 32, TRUE))
435 goto fail;
436
437 if (!InitializeCriticalSectionAndSpinCount(&pool->lock, 4000))
438 goto fail;
439
440 return pool;
441fail:
442 WINPR_PRAGMA_DIAG_PUSH
443 WINPR_PRAGMA_DIAG_IGNORED_MISMATCHED_DEALLOC
444 StreamPool_Free(pool);
445 WINPR_PRAGMA_DIAG_POP
446 return nullptr;
447}
448
449void StreamPool_Free(wStreamPool* pool)
450{
451 if (!pool)
452 return;
453
454 StreamPool_Clear(pool);
455
456 DeleteCriticalSection(&pool->lock);
457
458 free(pool->aArray);
459 free(pool->uArray);
460
461 WLog_Discard(pool->log);
462 free(pool);
463}
464
465char* StreamPool_GetStatistics(wStreamPool* pool, char* buffer, size_t size)
466{
467 WINPR_ASSERT(pool);
468
469 if (!buffer || (size < 1))
470 return nullptr;
471
472 size_t used = 0;
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;
479
480#if defined(WITH_DEBUG_STREAMPOOL)
481 StreamPool_Lock(pool);
482
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++)
487 {
488 const struct s_StreamPoolEntry* cur = &pool->uArray[x];
489 WINPR_ASSERT(cur->msg || (cur->lines == 0));
490
491 for (size_t y = 0; y < cur->lines; y++)
492 {
493 offset = _snprintf(&buffer[used], size - 1 - used, "[%" PRIuz " | %" PRIuz "]: %s\n", x,
494 y, cur->msg[y]);
495 if ((offset > 0) && ((size_t)offset < size - used))
496 used += (size_t)offset;
497 }
498 }
499
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;
503
504 struct s_StreamPoolEntry entry = WINPR_C_ARRAY_INIT;
505 void* stack = winpr_backtrace(20);
506 if (stack)
507 entry.msg = winpr_backtrace_symbols(stack, &entry.lines);
508 winpr_backtrace_free(stack);
509
510 for (size_t x = 0; x < entry.lines; x++)
511 {
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;
516 }
517 free((void*)entry.msg);
518 StreamPool_Unlock(pool);
519#endif
520 buffer[used] = '\0';
521 return buffer;
522}
523
524BOOL StreamPool_WaitForReturn(wStreamPool* pool, UINT32 timeoutMS)
525{
526 /* HACK: We disconnected the transport above, now wait without a read or write lock until all
527 * streams in use have been returned to the pool. */
528 while (timeoutMS > 0)
529 {
530 const size_t used = StreamPool_UsedCount(pool);
531 if (used == 0)
532 return TRUE;
533 WLog_Print(pool->log, WLOG_DEBUG, "%" PRIuz " streams still in use, sleeping...", used);
534
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);
538
539 UINT32 diff = 10;
540 if (timeoutMS != INFINITE)
541 {
542 diff = timeoutMS > 10 ? 10 : timeoutMS;
543 timeoutMS -= diff;
544 }
545 Sleep(diff);
546 }
547
548 return FALSE;
549}