@@ -1673,6 +1673,14 @@ def core_func(args):
16731673 lift_and_run (opts , inst , ft , core_func , on_start , on_resolve )
16741674 assert (host_writer .readable_end is return_value )
16751675
1676+ # A stream whose writable end has already been dropped can still be lowered
1677+ # into and lifted out of the component; whoever finally reads it sees DROPPED.
1678+ host_writer = HostWriter (U8Type (), [])
1679+ lift_and_run (opts , inst , ft , core_func , on_start , on_resolve )
1680+ assert (host_writer .readable_end is return_value )
1681+ host_reader = HostReader (return_value )
1682+ assert (host_reader .dropped and host_reader .take () == [])
1683+
16761684
16771685def test_receive_own_stream ():
16781686 store = Store ()
@@ -2254,6 +2262,66 @@ def core_func(args):
22542262 lift_and_run (opts , inst , caller_ft , core_func , lambda :[], lambda _ :())
22552263
22562264
2265+ def test_stream_drop_both_ends_while_idle ():
2266+ store = Store ()
2267+ inst = ComponentInstance (store )
2268+ mem = bytearray (24 )
2269+ opts = mk_opts (memory = MemInst (mem , 'i32' ), async_ = True )
2270+ stream_t = StreamType (U8Type ())
2271+ future_t = FutureType (U8Type ())
2272+
2273+ def core_func (args ):
2274+ assert (len (args ) == 0 )
2275+ [] = canon_task_return ([], opts , [])
2276+ [packed ] = canon_stream_new (stream_t )
2277+ rsi ,wsi = unpack_new_ends (packed )
2278+ [] = canon_stream_drop_writable (stream_t , wsi )
2279+ [] = canon_stream_drop_readable (stream_t , rsi )
2280+
2281+ # Dropping one end notifies the other end even when it has no read or write
2282+ # in flight: the DROPPED event is delivered (once) through the waitable set,
2283+ # after which the end is done and further reads trap.
2284+ retp = 8
2285+ [seti ] = canon_waitable_set_new ()
2286+ [packed ] = canon_stream_new (stream_t )
2287+ rsi ,wsi = unpack_new_ends (packed )
2288+ [] = canon_waitable_join (rsi , seti )
2289+ [] = canon_stream_drop_writable (stream_t , wsi )
2290+ [event ] = canon_waitable_set_poll (MemInst (mem , 'i32' ), seti , retp )
2291+ assert (event == EventCode .STREAM_READ )
2292+ assert (mem [retp + 0 ] == rsi )
2293+ result ,n = unpack_result (mem [retp + 4 ])
2294+ assert (n == 0 and result == End .Result .DROPPED )
2295+ [event ] = canon_waitable_set_poll (MemInst (mem , 'i32' ), seti , retp )
2296+ assert (event == EventCode .NONE )
2297+ trapped = False
2298+ try :
2299+ canon_stream_read (stream_t , opts , rsi , 0 , 4 )
2300+ except Trap :
2301+ trapped = True
2302+ assert (trapped )
2303+ [] = canon_waitable_join (rsi , 0 )
2304+ [] = canon_stream_drop_readable (stream_t , rsi )
2305+
2306+ # Same for the writable end of a future (using wait instead of poll); once
2307+ # notified, the writable end may be dropped without having written.
2308+ [packed ] = canon_future_new (future_t )
2309+ rfi ,wfi = unpack_new_ends (packed )
2310+ [] = canon_waitable_join (wfi , seti )
2311+ [] = canon_future_drop_readable (future_t , rfi )
2312+ [event ] = canon_waitable_set_wait (MemInst (mem , 'i32' ), seti , retp )
2313+ assert (event == EventCode .FUTURE_WRITE )
2314+ assert (mem [retp + 0 ] == wfi )
2315+ assert (mem [retp + 4 ] == End .Result .DROPPED )
2316+ [] = canon_waitable_join (wfi , 0 )
2317+ [] = canon_future_drop_writable (future_t , wfi )
2318+ [] = canon_waitable_set_drop (seti )
2319+ return []
2320+
2321+ caller_ft = FuncType ([], [], async_ = True )
2322+ lift_and_run (opts , inst , caller_ft , core_func , lambda :[], lambda _ :())
2323+
2324+
22572325def test_cancel_subtask ():
22582326 store = Store ()
22592327 ft = FuncType ([U8Type ()], [U8Type ()], async_ = True )
@@ -2982,6 +3050,7 @@ def on_resolve(v):
29823050test_cancel_copy ()
29833051test_futures ()
29843052test_future_drop_readable_with_pending_write ()
3053+ test_stream_drop_both_ends_while_idle ()
29853054test_cancel_subtask ()
29863055test_self_copy (None )
29873056test_self_copy (U8Type ())
0 commit comments