502
|
1 |
SequenceableCollection subclass:#CollectingSharedQueueStream
|
|
2 |
instanceVariableNames:'contents readPosition writePosition dataAvailable closed
|
|
3 |
finalSizeIfKnown signalChanges'
|
500
|
4 |
classVariableNames:''
|
|
5 |
poolDictionaries:''
|
501
|
6 |
category:'Streams'
|
500
|
7 |
!
|
|
8 |
|
|
9 |
!CollectingSharedQueueStream class methodsFor:'documentation'!
|
|
10 |
|
|
11 |
documentation
|
|
12 |
"
|
501
|
13 |
This class provides a buffering mechanism between a reader and a writer
|
|
14 |
process (i.e. it is much like a sharedQueue), but remembers the data as
|
|
15 |
written by the writer internally, providing indexed access to elements.
|
|
16 |
|
500
|
17 |
The reader side may read from it using #next, and possibly access
|
|
18 |
elements via #at:.
|
501
|
19 |
Reading/accessing may start immediately, but will block until enough elements
|
500
|
20 |
have been added by another process, the writer.
|
501
|
21 |
|
500
|
22 |
Instances of this class may be useful to start processing on
|
|
23 |
big document/data collection immediately, while the data is still being
|
501
|
24 |
read by another thread;
|
|
25 |
A concrete application is the HTMLDocumentReader, which is being changed to
|
|
26 |
start processing and displaying the document while the rest is still being read.
|
500
|
27 |
|
|
28 |
[author:]
|
|
29 |
Claus Gittinger
|
|
30 |
|
|
31 |
[see also:]
|
|
32 |
Stream OrderedCollection SharedQueue
|
|
33 |
"
|
|
34 |
!
|
|
35 |
|
|
36 |
examples
|
|
37 |
"
|
502
|
38 |
two processes synchronized much like with a sharedQueue:
|
500
|
39 |
[exBegin]
|
|
40 |
|s reader|
|
|
41 |
|
|
42 |
s := CollectingSharedQueueStream new.
|
|
43 |
reader := [
|
|
44 |
[s atEnd] whileFalse:[
|
|
45 |
Transcript showCR:s next
|
|
46 |
].
|
|
47 |
] fork.
|
|
48 |
|
|
49 |
1 to:10 do:[:i |
|
|
50 |
Delay waitForSeconds:1.
|
|
51 |
s nextPut:i
|
|
52 |
].
|
|
53 |
[exEnd]
|
|
54 |
|
502
|
55 |
the writer reads from a (slow) pipe;
|
|
56 |
the reader sends it to the transcript.
|
500
|
57 |
[exBegin]
|
|
58 |
|pipe s reader|
|
|
59 |
|
|
60 |
s := CollectingSharedQueueStream new.
|
|
61 |
reader := [
|
|
62 |
[s atEnd] whileFalse:[
|
|
63 |
Transcript showCR:s next
|
|
64 |
].
|
|
65 |
] forkAt:9.
|
|
66 |
|
|
67 |
pipe := PipeStream readingFrom:'ls -lR /usr'.
|
|
68 |
pipe notNil ifTrue:[
|
|
69 |
[pipe atEnd] whileFalse:[
|
|
70 |
pipe readWait.
|
|
71 |
s nextPut:(pipe nextLine).
|
|
72 |
].
|
|
73 |
pipe close.
|
|
74 |
].
|
|
75 |
s close
|
|
76 |
[exEnd]
|
502
|
77 |
|
|
78 |
|
|
79 |
the writer reads from a (slow) pipe;
|
|
80 |
the collection is used in a TextView, which
|
|
81 |
will block whenever lines are to be displayed, which have not
|
|
82 |
yet been read:
|
|
83 |
[exBegin]
|
|
84 |
|view pipe buffer reader|
|
|
85 |
|
|
86 |
buffer := CollectingSharedQueueStream new.
|
|
87 |
buffer finalSize:100.
|
|
88 |
|
|
89 |
[
|
|
90 |
pipe := PipeStream readingFrom:'ls -lR /usr'.
|
|
91 |
pipe notNil ifTrue:[
|
|
92 |
[pipe atEnd] whileFalse:[
|
|
93 |
pipe readWait.
|
|
94 |
buffer nextPut:(pipe nextLine).
|
|
95 |
].
|
|
96 |
pipe close.
|
|
97 |
].
|
|
98 |
buffer changed:#size.
|
|
99 |
buffer close.
|
|
100 |
] fork.
|
|
101 |
|
|
102 |
view := ScrollableView for:TextView.
|
|
103 |
view model:buffer; listMessage:#value; aspectMessage:#value.
|
|
104 |
view open.
|
|
105 |
[exEnd]
|
500
|
106 |
"
|
|
107 |
|
|
108 |
! !
|
|
109 |
|
|
110 |
!CollectingSharedQueueStream class methodsFor:'instance creation'!
|
|
111 |
|
|
112 |
new
|
|
113 |
^ self basicNew initialize
|
|
114 |
|
|
115 |
"Created: 5.3.1997 / 14:30:36 / cg"
|
|
116 |
! !
|
|
117 |
|
|
118 |
!CollectingSharedQueueStream methodsFor:'accessing'!
|
|
119 |
|
|
120 |
at:index
|
|
121 |
"synchronized read - possibly wait for elements up to index
|
|
122 |
being added (by someone else); then return it."
|
|
123 |
|
|
124 |
writePosition > index ifTrue:[
|
|
125 |
^ contents at:index
|
|
126 |
].
|
|
127 |
|
|
128 |
[writePosition <= index] whileTrue:[
|
|
129 |
closed ifTrue:[
|
|
130 |
^ self subscriptBoundsError:index
|
|
131 |
].
|
|
132 |
dataAvailable wait.
|
|
133 |
].
|
|
134 |
^ contents at:index
|
|
135 |
|
|
136 |
"Created: 5.3.1997 / 14:44:41 / cg"
|
|
137 |
!
|
|
138 |
|
|
139 |
next
|
|
140 |
"return the next value in the queue; if there is none,
|
|
141 |
wait 'til something is put into the receiver."
|
|
142 |
|
|
143 |
|value|
|
|
144 |
|
|
145 |
[readPosition >= writePosition] whileTrue:[
|
|
146 |
closed ifTrue:[
|
|
147 |
^ nil
|
|
148 |
].
|
|
149 |
dataAvailable wait
|
|
150 |
].
|
|
151 |
|
|
152 |
value := contents at:readPosition.
|
|
153 |
readPosition := readPosition + 1.
|
|
154 |
^ value
|
|
155 |
|
|
156 |
"Created: 5.3.1997 / 14:28:57 / cg"
|
|
157 |
"Modified: 5.3.1997 / 14:45:54 / cg"
|
|
158 |
!
|
|
159 |
|
|
160 |
nextPut:anObject
|
|
161 |
"append anObject to the queue; if anyone is waiting, tell him"
|
|
162 |
|
|
163 |
|value|
|
|
164 |
|
|
165 |
contents add:anObject.
|
|
166 |
writePosition := writePosition + 1.
|
502
|
167 |
|
|
168 |
finalSizeIfKnown isNil ifTrue:[
|
|
169 |
signalChanges ifTrue:[
|
|
170 |
self changed:#size with:nil
|
|
171 |
]
|
|
172 |
].
|
500
|
173 |
dataAvailable signal
|
|
174 |
|
|
175 |
"Created: 5.3.1997 / 14:33:44 / cg"
|
502
|
176 |
"Modified: 5.3.1997 / 15:42:33 / cg"
|
|
177 |
! !
|
|
178 |
|
|
179 |
!CollectingSharedQueueStream methodsFor:'accessing - special'!
|
|
180 |
|
|
181 |
close
|
|
182 |
"signal the end of input; to be used by the writer"
|
|
183 |
|
|
184 |
closed := true.
|
|
185 |
dataAvailable signal
|
|
186 |
|
|
187 |
"Modified: 5.3.1997 / 14:45:11 / cg"
|
|
188 |
!
|
|
189 |
|
|
190 |
finalSize:aNumber
|
|
191 |
"can be used by the writer, if the final size is known in
|
|
192 |
advance."
|
|
193 |
|
|
194 |
finalSizeIfKnown := aNumber.
|
|
195 |
signalChanges ifTrue:[
|
|
196 |
self changed:#size
|
|
197 |
].
|
|
198 |
|
|
199 |
"Created: 5.3.1997 / 15:36:07 / cg"
|
|
200 |
"Modified: 5.3.1997 / 15:57:24 / cg"
|
|
201 |
!
|
|
202 |
|
|
203 |
signalChanges:aBoolean
|
|
204 |
"controls if I should send out size-changeMessages when new elements arrive"
|
|
205 |
|
|
206 |
signalChanges := aBoolean.
|
|
207 |
|
|
208 |
"Created: 5.3.1997 / 15:40:57 / cg"
|
|
209 |
! !
|
|
210 |
|
|
211 |
!CollectingSharedQueueStream methodsFor:'dummy converting'!
|
|
212 |
|
|
213 |
asStringCollection
|
|
214 |
^ self
|
|
215 |
|
|
216 |
"Created: 5.3.1997 / 16:02:57 / cg"
|
500
|
217 |
! !
|
|
218 |
|
|
219 |
!CollectingSharedQueueStream methodsFor:'initialization'!
|
|
220 |
|
|
221 |
initialize
|
|
222 |
readPosition := writePosition := 1.
|
|
223 |
dataAvailable := Semaphore new.
|
|
224 |
contents := OrderedCollection new.
|
|
225 |
closed := false.
|
|
226 |
|
|
227 |
"Modified: 5.3.1997 / 14:34:55 / cg"
|
|
228 |
! !
|
|
229 |
|
|
230 |
!CollectingSharedQueueStream methodsFor:'queries'!
|
|
231 |
|
|
232 |
atEnd
|
|
233 |
closed ifFalse:[^ false].
|
|
234 |
^ readPosition >= writePosition
|
|
235 |
|
|
236 |
"Modified: 5.3.1997 / 14:41:04 / cg"
|
502
|
237 |
!
|
|
238 |
|
|
239 |
currentSize
|
|
240 |
^ contents size
|
|
241 |
|
|
242 |
"Modified: 5.3.1997 / 15:56:36 / cg"
|
|
243 |
!
|
|
244 |
|
|
245 |
size
|
|
246 |
closed ifTrue:[^ contents size].
|
|
247 |
finalSizeIfKnown notNil ifTrue:[^ finalSizeIfKnown].
|
|
248 |
|
|
249 |
"/ must wait until closed
|
|
250 |
[closed] whileFalse:[
|
|
251 |
dataAvailable wait
|
|
252 |
].
|
|
253 |
closed ifTrue:[^ contents size].
|
|
254 |
|
|
255 |
"Created: 5.3.1997 / 15:35:29 / cg"
|
|
256 |
"Modified: 5.3.1997 / 15:57:08 / cg"
|
500
|
257 |
! !
|
|
258 |
|
|
259 |
!CollectingSharedQueueStream class methodsFor:'documentation'!
|
|
260 |
|
|
261 |
version
|
502
|
262 |
^ '$Header: /cvs/stx/stx/libbasic2/CollectingSharedQueueStream.st,v 1.3 1997-03-05 16:26:06 cg Exp $'
|
500
|
263 |
! !
|