forked from SerenityOS/serenity
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathAsyncStreamTransform.h
115 lines (97 loc) Β· 2.95 KB
/
AsyncStreamTransform.h
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
/*
* Copyright (c) 2024, Dan Klishch <[email protected]>
*
* SPDX-License-Identifier: BSD-2-Clause
*/
#pragma once
#include <AK/AsyncStream.h>
#include <AK/Generator.h>
#include <AK/MaybeOwned.h>
#include <AK/TemporaryChange.h>
namespace AK {
template<typename T>
class AsyncStreamTransform : public AsyncInputStream {
public:
AsyncStreamTransform(MaybeOwned<T>&& stream, AK::Generator<Empty, ErrorOr<void>>&& generator)
: m_stream(move(stream))
, m_generator(move(generator))
{
}
~AsyncStreamTransform()
{
// 1. Assert that nobody is awaiting on the resource.
VERIFY(!m_generator_has_awaiters);
// 2. If resource is open, perform Reset AO.
if (is_open())
reset();
}
void reset() override
{
VERIFY(is_open());
m_stream->reset();
if (!m_generator_has_awaiters)
m_generator.destroy();
m_is_open = false;
}
Coroutine<ErrorOr<void>> close() override
{
VERIFY(is_open());
TemporaryChange await_guard(m_generator_has_awaiters, true);
if (!m_generator.is_done()) {
Variant<Empty, ErrorOr<void>> chunk_or_eof = co_await m_generator.next();
if (chunk_or_eof.has<Empty>()) {
reset();
co_return Error::from_errno(EBUSY);
} else {
m_is_open = false;
auto& error_or_eof = chunk_or_eof.get<ErrorOr<void>>();
if (error_or_eof.is_error()) {
co_return error_or_eof.release_error();
} else {
if (m_stream.is_owned())
CO_TRY(co_await m_stream->close());
co_return {};
}
}
} else {
m_is_open = false;
if (m_stream.is_owned())
CO_TRY(co_await m_stream->close());
co_return {};
}
}
bool is_open() const override
{
return m_is_open;
}
Coroutine<ErrorOr<bool>> enqueue_some(Badge<AsyncInputStream>) override
{
VERIFY(is_open());
TemporaryChange await_guard(m_generator_has_awaiters, true);
if (m_generator.is_done())
co_return false;
Variant<Empty, ErrorOr<void>> chunk_or_eof = co_await m_generator.next();
if (chunk_or_eof.has<Empty>()) {
co_return true;
} else {
auto& error_or_eof = chunk_or_eof.get<ErrorOr<void>>();
if (error_or_eof.is_error()) {
m_is_open = false;
co_return error_or_eof.release_error();
} else {
co_return false;
}
}
}
protected:
using Generator = AK::Generator<Empty, ErrorOr<void>>;
MaybeOwned<T> m_stream;
private:
Generator m_generator;
bool m_is_open { true };
bool m_generator_has_awaiters { false };
};
}
#ifdef USING_AK_GLOBALLY
using AK::AsyncStreamTransform;
#endif