Isolate cached feed broadcast consumers #92
|
|
@ -0,0 +1,34 @@
|
||||||
|
---
|
||||||
|
model: openai/gpt-5.6-sol
|
||||||
|
service: opencode
|
||||||
|
session: unavailable
|
||||||
|
timestamp: 2026-08-13T20:24:02Z
|
||||||
|
git_ref: 12b0d66e
|
||||||
|
scope: code
|
||||||
|
substantive: true
|
||||||
|
raw_file: 20260813T202402Z_12b0d66e_prompt_io.raw.md
|
||||||
|
---
|
||||||
|
|
||||||
|
## Prompt
|
||||||
|
|
||||||
|
Continue the paired piker broadcast-consumer work after Tractor PR 490
|
||||||
|
stabilization. Stop cached feed callers from mutating or consuming shared
|
||||||
|
quote-stream roots, then review, verify and prepare a commit boundary.
|
||||||
|
|
||||||
|
## Response summary
|
||||||
|
|
||||||
|
Made provider root streams private to the cached `Feed`, gave every caller a
|
||||||
|
local subscribed view, corrected feed cache identity, and made multi-provider
|
||||||
|
relays own synchronized child subscriptions and independent send-channel
|
||||||
|
clones.
|
||||||
|
|
||||||
|
## Files changed
|
||||||
|
|
||||||
|
- `piker/data/feed.py` - caller-local feed views, complete cache keys, and
|
||||||
|
subscribed multi-provider relays with independent channel lifetimes.
|
||||||
|
- `tests/test_feed_broadcast.py` - real broadcast ownership, lifecycle and
|
||||||
|
cache-identity regressions, including sibling survival after provider EOC.
|
||||||
|
|
||||||
|
## Human edits
|
||||||
|
|
||||||
|
None - generated output follows the requested cached-feed ownership item.
|
||||||
|
|
@ -0,0 +1,51 @@
|
||||||
|
---
|
||||||
|
model: openai/gpt-5.6-sol
|
||||||
|
service: opencode
|
||||||
|
timestamp: 2026-08-13T20:24:02Z
|
||||||
|
git_ref: 12b0d66e
|
||||||
|
diff_cmd: git diff HEAD~1..HEAD
|
||||||
|
---
|
||||||
|
|
||||||
|
The user asked to continue the paired piker broadcast-consumer work after
|
||||||
|
Tractor PR 490 stabilization. The next item was to stop cached feed callers
|
||||||
|
from mutating or consuming shared `Feed`/`Flume` quote-stream roots, then
|
||||||
|
review, verify and prepare a commit boundary.
|
||||||
|
|
||||||
|
> `git diff HEAD~1..HEAD -- piker/data/feed.py`
|
||||||
|
|
||||||
|
Changed `maybe_open_feed()` so the cached `Feed` exclusively owns provider
|
||||||
|
root streams and every caller, including the first cache owner, receives a
|
||||||
|
caller-local `Feed`/`Flume` view backed by one child subscription per provider.
|
||||||
|
The cache identity now includes normalized full symbol shape and stream policy.
|
||||||
|
|
||||||
|
Changed `Feed.open_multi_stream()` so every provider relay subscribes instead
|
||||||
|
of directly consuming a shared stream, and synchronizes relay startup before
|
||||||
|
yielding the merged channel. Each relay owns a cloned send channel so one
|
||||||
|
provider EOC cannot terminate its siblings. The module-level relay returns its
|
||||||
|
live child receiver through `TaskStatus` for typed retention and introspection.
|
||||||
|
Documented `Feed.pause()`/`Feed.resume()` as provider-context-wide controls
|
||||||
|
shared by local cache users.
|
||||||
|
|
||||||
|
> `git diff HEAD~1..HEAD -- tests/test_feed_broadcast.py`
|
||||||
|
|
||||||
|
Added real `Feed`, `Flume`, and Tractor broadcast regressions covering both
|
||||||
|
cache-miss owner and cache-hit paths, one-provider-many-market mapping,
|
||||||
|
multi-provider relay fan-out and startup, duplex controls, lexical child
|
||||||
|
closure, sibling survival after provider EOC, cached-root immutability, and
|
||||||
|
full cache-key differentiation.
|
||||||
|
|
||||||
|
Verification output:
|
||||||
|
|
||||||
|
```text
|
||||||
|
... [100%]
|
||||||
|
3 passed in 0.02s
|
||||||
|
|
||||||
|
. [100%]
|
||||||
|
1 passed in 0.01s
|
||||||
|
```
|
||||||
|
|
||||||
|
The second result is the existing EMS broadcast ownership regression. Ruff,
|
||||||
|
E501, compilation and whitespace checks passed. Repeated adversarial review
|
||||||
|
found and resolved shared `Feed.streams`, incomplete cache-key, fake-lifecycle,
|
||||||
|
relay-startup and first-owner asymmetry issues; final review reported no
|
||||||
|
findings.
|
||||||
|
|
@ -0,0 +1,32 @@
|
||||||
|
---
|
||||||
|
model: openai/gpt-5.6-sol
|
||||||
|
service: opencode
|
||||||
|
session: unavailable
|
||||||
|
timestamp: 2026-08-13T22:59:40Z
|
||||||
|
git_ref: b691422b
|
||||||
|
scope: code
|
||||||
|
substantive: true
|
||||||
|
raw_file: 20260813T225940Z_b691422b_prompt_io.raw.md
|
||||||
|
---
|
||||||
|
|
||||||
|
## Prompt
|
||||||
|
|
||||||
|
Continue the piker broadcast-consumer audit after the cached feed ownership
|
||||||
|
commit and prepare the next smallest ownership fix.
|
||||||
|
|
||||||
|
## Response summary
|
||||||
|
|
||||||
|
Made every Deribit cached price-feed caller own a lexical child subscription
|
||||||
|
from the retained `LinkedTaskChannel` broadcaster instead of constructing
|
||||||
|
competing roots around one source channel.
|
||||||
|
|
||||||
|
## Files changed
|
||||||
|
|
||||||
|
- `piker/brokers/deribit/api.py` - uniform child subscription ownership for
|
||||||
|
price-feed misses and cache hits.
|
||||||
|
- `tests/test_deribit_broadcast.py` - concurrent fan-out and child lifetime
|
||||||
|
regression.
|
||||||
|
|
||||||
|
## Human edits
|
||||||
|
|
||||||
|
None - generated output follows the requested broadcast-consumer audit.
|
||||||
|
|
@ -0,0 +1,37 @@
|
||||||
|
---
|
||||||
|
model: openai/gpt-5.6-sol
|
||||||
|
service: opencode
|
||||||
|
timestamp: 2026-08-13T22:59:40Z
|
||||||
|
git_ref: b691422b
|
||||||
|
diff_cmd: git diff HEAD~1..HEAD
|
||||||
|
---
|
||||||
|
|
||||||
|
The user asked to continue the piker broadcast-consumer audit after committing
|
||||||
|
the cached feed ownership fix. The next isolated item was Deribit's cached
|
||||||
|
price feed, which wrapped one retained `LinkedTaskChannel` in competing
|
||||||
|
broadcaster roots for cache hits.
|
||||||
|
|
||||||
|
> `git diff HEAD~1..HEAD -- piker/brokers/deribit/api.py`
|
||||||
|
|
||||||
|
Changed `maybe_open_price_feed()` so both the first cache owner and every cache
|
||||||
|
hit enter `LinkedTaskChannel.subscribe()` and receive a lexical child from the
|
||||||
|
channel's single retained broadcaster root.
|
||||||
|
|
||||||
|
> `git diff HEAD~1..HEAD -- tests/test_deribit_broadcast.py`
|
||||||
|
|
||||||
|
Added a concurrent ownership regression using real Tractor broadcast state.
|
||||||
|
It proves a miss owner and cache-hit caller receive the same quote, closing the
|
||||||
|
hit child leaves the owner live, and both children clean up independently while
|
||||||
|
the channel root remains registered.
|
||||||
|
|
||||||
|
Verification output:
|
||||||
|
|
||||||
|
```text
|
||||||
|
. [100%]
|
||||||
|
1 passed in 0.01s
|
||||||
|
```
|
||||||
|
|
||||||
|
The new test passes through a source-isolated Deribit import. Test Ruff,
|
||||||
|
Python compilation and whitespace checks pass. Full package import and
|
||||||
|
full-file Deribit Ruff remain blocked by pre-existing Qt, msgspec model, and
|
||||||
|
legacy Deribit lint errors unrelated to this patch.
|
||||||
|
|
@ -0,0 +1,32 @@
|
||||||
|
---
|
||||||
|
model: openai/gpt-5.6-sol
|
||||||
|
service: opencode
|
||||||
|
session: unavailable
|
||||||
|
timestamp: 2026-08-13T23:07:44Z
|
||||||
|
git_ref: ddac3553
|
||||||
|
scope: code
|
||||||
|
substantive: true
|
||||||
|
raw_file: 20260813T230744Z_ddac3553_prompt_io.raw.md
|
||||||
|
---
|
||||||
|
|
||||||
|
## Prompt
|
||||||
|
|
||||||
|
Continue the piker broadcast-consumer audit after the Deribit ownership commit
|
||||||
|
and prepare the next smallest ownership fix.
|
||||||
|
|
||||||
|
## Response summary
|
||||||
|
|
||||||
|
Replaced IB's manual raw quote-channel registry with a symbol-keyed cached
|
||||||
|
context and gave every caller a lexical child from one retained
|
||||||
|
`LinkedTaskChannel` broadcaster root.
|
||||||
|
|
||||||
|
## Files changed
|
||||||
|
|
||||||
|
- `piker/brokers/ib/feed.py` - cached raw-channel lifetime and lexical child
|
||||||
|
ownership.
|
||||||
|
- `tests/test_ib_broadcast.py` - concurrent fan-out, cache-key and child
|
||||||
|
lifetime regression.
|
||||||
|
|
||||||
|
## Human edits
|
||||||
|
|
||||||
|
None - generated output follows the requested broadcast-consumer audit.
|
||||||
|
|
@ -0,0 +1,39 @@
|
||||||
|
---
|
||||||
|
model: openai/gpt-5.6-sol
|
||||||
|
service: opencode
|
||||||
|
timestamp: 2026-08-13T23:07:44Z
|
||||||
|
git_ref: ddac3553
|
||||||
|
diff_cmd: git diff HEAD~1..HEAD
|
||||||
|
---
|
||||||
|
|
||||||
|
The user asked to continue the piker broadcast-consumer audit after committing
|
||||||
|
the Deribit ownership fix. The next isolated item was IB's actor-local
|
||||||
|
`_quote_streams` cache, which retained a raw channel owned by its first caller
|
||||||
|
and constructed competing broadcaster roots for later callers.
|
||||||
|
|
||||||
|
> `git diff HEAD~1..HEAD -- piker/brokers/ib/feed.py`
|
||||||
|
|
||||||
|
Factored `_open_aio_quote_channel()` as the raw asyncio-channel context and
|
||||||
|
made `open_aio_quote_stream()` cache that context by symbol through
|
||||||
|
`maybe_open_context()`. Every caller now receives a lexical child from the
|
||||||
|
retained `LinkedTaskChannel` broadcaster, while the raw source remains alive
|
||||||
|
until its final cache user exits. Removed the manual `_quote_streams` registry
|
||||||
|
and teardown mutation.
|
||||||
|
|
||||||
|
> `git diff HEAD~1..HEAD -- tests/test_ib_broadcast.py`
|
||||||
|
|
||||||
|
Added a concurrent ownership regression using real Tractor broadcast state.
|
||||||
|
It proves miss and hit callers use the same symbol/context factory, receive the
|
||||||
|
same quote through distinct children, close independently, and leave the owner
|
||||||
|
live after the hit exits.
|
||||||
|
|
||||||
|
Verification output:
|
||||||
|
|
||||||
|
```text
|
||||||
|
. [100%]
|
||||||
|
1 passed in 0.01s
|
||||||
|
```
|
||||||
|
|
||||||
|
The new test passes through a source-isolated IB feed import. Test Ruff, E501,
|
||||||
|
Python compilation and whitespace checks pass. Normal package collection
|
||||||
|
remains blocked by the existing Qt environment mismatch.
|
||||||
|
|
@ -36,6 +36,7 @@ from trio_typing import TaskStatus
|
||||||
from rapidfuzz import process as fuzzy
|
from rapidfuzz import process as fuzzy
|
||||||
import numpy as np
|
import numpy as np
|
||||||
from tractor.trionics import (
|
from tractor.trionics import (
|
||||||
|
BroadcastReceiver,
|
||||||
broadcast_receiver,
|
broadcast_receiver,
|
||||||
maybe_open_context,
|
maybe_open_context,
|
||||||
collapse_eg,
|
collapse_eg,
|
||||||
|
|
@ -589,20 +590,20 @@ async def open_price_feed(
|
||||||
@acm
|
@acm
|
||||||
async def maybe_open_price_feed(
|
async def maybe_open_price_feed(
|
||||||
instrument: str
|
instrument: str
|
||||||
) -> trio.abc.ReceiveStream:
|
) -> BroadcastReceiver:
|
||||||
|
|
||||||
# TODO: add a predicate to maybe_open_context
|
# TODO: add a predicate to maybe_open_context
|
||||||
|
feed: to_asyncio.LinkedTaskChannel
|
||||||
async with maybe_open_context(
|
async with maybe_open_context(
|
||||||
acm_func=open_price_feed,
|
acm_func=open_price_feed,
|
||||||
kwargs={
|
kwargs={
|
||||||
'instrument': instrument
|
'instrument': instrument
|
||||||
},
|
},
|
||||||
key=f'{instrument}-price',
|
key=f'{instrument}-price',
|
||||||
) as (cache_hit, feed):
|
) as (_, feed):
|
||||||
if cache_hit:
|
bstream: BroadcastReceiver
|
||||||
yield broadcast_receiver(feed, 10)
|
async with feed.subscribe() as bstream:
|
||||||
else:
|
yield bstream
|
||||||
yield feed
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -738,10 +738,6 @@ async def get_bars(
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
# per-actor cache of inter-eventloop-chans
|
|
||||||
_quote_streams: dict[str, trio.abc.ReceiveStream] = {}
|
|
||||||
|
|
||||||
|
|
||||||
async def _setup_quote_stream(
|
async def _setup_quote_stream(
|
||||||
chan: tractor.to_asyncio.LinkedTaskChannel,
|
chan: tractor.to_asyncio.LinkedTaskChannel,
|
||||||
symbol: str,
|
symbol: str,
|
||||||
|
|
@ -770,8 +766,6 @@ async def _setup_quote_stream(
|
||||||
and is thus run via `tractor.to_asyncio.open_channel_from()`.
|
and is thus run via `tractor.to_asyncio.open_channel_from()`.
|
||||||
|
|
||||||
'''
|
'''
|
||||||
global _quote_streams
|
|
||||||
|
|
||||||
async with load_aio_clients(
|
async with load_aio_clients(
|
||||||
disconnect_on_exit=False,
|
disconnect_on_exit=False,
|
||||||
) as accts2clients:
|
) as accts2clients:
|
||||||
|
|
@ -834,9 +828,6 @@ async def _setup_quote_stream(
|
||||||
|
|
||||||
client.ib.cancelMktData(contract)
|
client.ib.cancelMktData(contract)
|
||||||
|
|
||||||
# decouple broadcast mem chan
|
|
||||||
_quote_streams.pop(symbol, None)
|
|
||||||
|
|
||||||
def push(
|
def push(
|
||||||
t: Ticker,
|
t: Ticker,
|
||||||
tries_before_raise: int = 6,
|
tries_before_raise: int = 6,
|
||||||
|
|
@ -953,54 +944,52 @@ async def _setup_quote_stream(
|
||||||
|
|
||||||
|
|
||||||
@acm
|
@acm
|
||||||
async def open_aio_quote_stream(
|
async def _open_aio_quote_channel(
|
||||||
symbol: str,
|
symbol: str,
|
||||||
contract: Contract|None = None,
|
contract: Contract|None = None,
|
||||||
|
|
||||||
) -> (
|
) -> tractor.to_asyncio.LinkedTaskChannel:
|
||||||
trio.abc.Channel| # iface
|
|
||||||
tractor.to_asyncio.LinkedTaskChannel # actually
|
|
||||||
):
|
|
||||||
'''
|
'''
|
||||||
Open a real-time `Ticker` quote stream from an `asyncio.Task`
|
Open one raw inter-event-loop quote channel for `symbol`.
|
||||||
spawned via `tractor.to_asyncio.open_channel_from()`, deliver the
|
|
||||||
inter-event-loop channel to the `trio.Task` caller and cache it
|
|
||||||
globally for re-use.
|
|
||||||
|
|
||||||
'''
|
'''
|
||||||
from tractor.trionics import broadcast_receiver
|
chan: tractor.to_asyncio.LinkedTaskChannel
|
||||||
global _quote_streams
|
started_contract: Contract
|
||||||
|
|
||||||
from_aio = _quote_streams.get(symbol)
|
|
||||||
if from_aio:
|
|
||||||
|
|
||||||
# if we already have a cached feed deliver a rx side clone
|
|
||||||
# to consumer
|
|
||||||
async with broadcast_receiver(
|
|
||||||
from_aio,
|
|
||||||
2**6,
|
|
||||||
) as from_aio:
|
|
||||||
yield from_aio
|
|
||||||
return
|
|
||||||
|
|
||||||
from_aio: tractor.to_asyncio.LinkedTaskChannel
|
|
||||||
async with tractor.to_asyncio.open_channel_from(
|
async with tractor.to_asyncio.open_channel_from(
|
||||||
_setup_quote_stream,
|
_setup_quote_stream,
|
||||||
symbol=symbol,
|
symbol=symbol,
|
||||||
contract=contract,
|
contract=contract,
|
||||||
|
|
||||||
) as (from_aio, contract):
|
) as (chan, started_contract):
|
||||||
|
assert started_contract
|
||||||
|
yield chan
|
||||||
|
|
||||||
assert contract
|
|
||||||
|
|
||||||
# TODO? de-reg on teardown of last consumer task?
|
@acm
|
||||||
# -> why aren't we using `.trionics.maybe_open_context()`
|
async def open_aio_quote_stream(
|
||||||
# here again?? (we are in `open_client_proxies()` tho?)
|
symbol: str,
|
||||||
#
|
contract: Contract|None = None,
|
||||||
# cache feed for later consumers
|
|
||||||
_quote_streams[symbol] = from_aio
|
|
||||||
|
|
||||||
yield from_aio
|
) -> tractor.trionics.BroadcastReceiver:
|
||||||
|
'''
|
||||||
|
Open a caller-owned child of the cached quote channel for `symbol`.
|
||||||
|
|
||||||
|
Keep the raw `LinkedTaskChannel` alive until its last cache user
|
||||||
|
exits, and retain exactly one broadcaster root on that channel.
|
||||||
|
|
||||||
|
'''
|
||||||
|
chan: tractor.to_asyncio.LinkedTaskChannel
|
||||||
|
async with tractor.trionics.maybe_open_context(
|
||||||
|
acm_func=_open_aio_quote_channel,
|
||||||
|
kwargs={
|
||||||
|
'symbol': symbol,
|
||||||
|
'contract': contract,
|
||||||
|
},
|
||||||
|
key=symbol,
|
||||||
|
) as (_, chan):
|
||||||
|
bstream: tractor.trionics.BroadcastReceiver
|
||||||
|
async with chan.subscribe() as bstream:
|
||||||
|
yield bstream
|
||||||
|
|
||||||
|
|
||||||
# TODO: cython/mypyc/numba this!
|
# TODO: cython/mypyc/numba this!
|
||||||
|
|
|
||||||
|
|
@ -159,11 +159,11 @@ class DarkBook(Struct):
|
||||||
_DEFAULT_SIZE: float = 1.0
|
_DEFAULT_SIZE: float = 1.0
|
||||||
|
|
||||||
|
|
||||||
async def clear_dark_triggers(
|
async def _clear_dark_triggers(
|
||||||
|
|
||||||
router: Router,
|
router: Router,
|
||||||
brokerd_orders_stream: tractor.MsgStream,
|
brokerd_orders_stream: tractor.MsgStream,
|
||||||
quote_stream: tractor.MsgStream,
|
quote_stream: trionics.AsyncReceiver,
|
||||||
broker: str,
|
broker: str,
|
||||||
fqme: str,
|
fqme: str,
|
||||||
|
|
||||||
|
|
@ -182,7 +182,6 @@ async def clear_dark_triggers(
|
||||||
# - port to the new ringbuf stuff in `tractor.ipc`!
|
# - port to the new ringbuf stuff in `tractor.ipc`!
|
||||||
# - numba all this!
|
# - numba all this!
|
||||||
# - this stream may eventually contain multiple symbols
|
# - this stream may eventually contain multiple symbols
|
||||||
quote_stream._raise_on_lag = False
|
|
||||||
async for quotes in quote_stream:
|
async for quotes in quote_stream:
|
||||||
# start = time.time()
|
# start = time.time()
|
||||||
for sym, quote in quotes.items():
|
for sym, quote in quotes.items():
|
||||||
|
|
@ -314,6 +313,32 @@ async def clear_dark_triggers(
|
||||||
# print(f'execs scan took: {time.time() - start}')
|
# print(f'execs scan took: {time.time() - start}')
|
||||||
|
|
||||||
|
|
||||||
|
async def clear_dark_triggers(
|
||||||
|
router: Router,
|
||||||
|
brokerd_orders_stream: tractor.MsgStream,
|
||||||
|
quote_stream: tractor.MsgStream,
|
||||||
|
broker: str,
|
||||||
|
fqme: str,
|
||||||
|
book: DarkBook,
|
||||||
|
|
||||||
|
) -> None:
|
||||||
|
'''
|
||||||
|
Run dark clearing on a dedicated non-strict quote subscription.
|
||||||
|
|
||||||
|
'''
|
||||||
|
async with quote_stream.subscribe(
|
||||||
|
raise_on_lag=False,
|
||||||
|
) as quotes_stream:
|
||||||
|
await _clear_dark_triggers(
|
||||||
|
router,
|
||||||
|
brokerd_orders_stream,
|
||||||
|
quotes_stream,
|
||||||
|
broker,
|
||||||
|
fqme,
|
||||||
|
book,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
class TradesRelay(Struct):
|
class TradesRelay(Struct):
|
||||||
|
|
||||||
# for now we keep only a single connection open with
|
# for now we keep only a single connection open with
|
||||||
|
|
|
||||||
|
|
@ -30,6 +30,7 @@ from collections import (
|
||||||
defaultdict,
|
defaultdict,
|
||||||
abc,
|
abc,
|
||||||
)
|
)
|
||||||
|
from copy import copy
|
||||||
from contextlib import asynccontextmanager as acm
|
from contextlib import asynccontextmanager as acm
|
||||||
from functools import partial
|
from functools import partial
|
||||||
import time
|
import time
|
||||||
|
|
@ -677,6 +678,36 @@ async def open_feed_bus(
|
||||||
bus.remove_subs(bs_fqme, subs)
|
bus.remove_subs(bs_fqme, subs)
|
||||||
|
|
||||||
|
|
||||||
|
async def relay_to_common_memchan(
|
||||||
|
stream: (
|
||||||
|
tractor.MsgStream
|
||||||
|
|trionics.BroadcastReceiver
|
||||||
|
),
|
||||||
|
tx: trio.MemorySendChannel[dict[str, Any]],
|
||||||
|
task_status: TaskStatus[
|
||||||
|
trionics.BroadcastReceiver,
|
||||||
|
] = trio.TASK_STATUS_IGNORED,
|
||||||
|
|
||||||
|
) -> None:
|
||||||
|
'''
|
||||||
|
Relay one provider subscription into a common memory channel.
|
||||||
|
|
||||||
|
Deliver the live child receiver through `task_status` so the
|
||||||
|
caller can retain and inspect each provider relay's cursor.
|
||||||
|
|
||||||
|
'''
|
||||||
|
bstream: trionics.BroadcastReceiver
|
||||||
|
async with (
|
||||||
|
stream.subscribe() as bstream,
|
||||||
|
tx,
|
||||||
|
):
|
||||||
|
task_status.started(bstream)
|
||||||
|
|
||||||
|
msg: dict[str, Any]
|
||||||
|
async for msg in bstream:
|
||||||
|
await tx.send(msg)
|
||||||
|
|
||||||
|
|
||||||
class Feed(Struct):
|
class Feed(Struct):
|
||||||
'''
|
'''
|
||||||
A per-provider API for client-side consumption from real-time data
|
A per-provider API for client-side consumption from real-time data
|
||||||
|
|
@ -726,29 +757,51 @@ class Feed(Struct):
|
||||||
if len(mods) == 1:
|
if len(mods) == 1:
|
||||||
# just pass the datad stream directly if only one provider
|
# just pass the datad stream directly if only one provider
|
||||||
# was detected.
|
# was detected.
|
||||||
stream = self.streams[list(brokers)[0]]
|
stream: (
|
||||||
|
tractor.MsgStream
|
||||||
|
|trionics.BroadcastReceiver
|
||||||
|
) = self.streams[list(brokers)[0]]
|
||||||
|
bstream: trionics.BroadcastReceiver
|
||||||
async with stream.subscribe() as bstream:
|
async with stream.subscribe() as bstream:
|
||||||
yield bstream
|
yield bstream
|
||||||
return
|
return
|
||||||
|
|
||||||
# start multiplexing task tree
|
# start multiplexing task tree
|
||||||
|
tx: trio.MemorySendChannel[dict[str, Any]]
|
||||||
|
rx: trio.MemoryReceiveChannel[dict[str, Any]]
|
||||||
tx, rx = trio.open_memory_channel(616)
|
tx, rx = trio.open_memory_channel(616)
|
||||||
|
|
||||||
async def relay_to_common_memchan(stream: tractor.MsgStream):
|
|
||||||
async with tx:
|
|
||||||
async for msg in stream:
|
|
||||||
await tx.send(msg)
|
|
||||||
|
|
||||||
async with (
|
async with (
|
||||||
tractor.trionics.collapse_eg(),
|
tractor.trionics.collapse_eg(),
|
||||||
trio.open_nursery() as nurse
|
trio.open_nursery() as nurse
|
||||||
):
|
):
|
||||||
# spawn a relay task for each stream so that they all
|
# spawn a relay task for each stream so that they all
|
||||||
# multiplex to a common channel.
|
# multiplex to a common channel.
|
||||||
for brokername in mods:
|
broker2bstreams: dict[
|
||||||
stream = self.streams[brokername]
|
str,
|
||||||
nurse.start_soon(relay_to_common_memchan, stream)
|
trionics.BroadcastReceiver,
|
||||||
|
] = {}
|
||||||
|
broker: str
|
||||||
|
for broker in mods:
|
||||||
|
stream: (
|
||||||
|
tractor.MsgStream
|
||||||
|
|trionics.BroadcastReceiver
|
||||||
|
) = self.streams[broker]
|
||||||
|
|
||||||
|
relay_tx: trio.MemorySendChannel[
|
||||||
|
dict[str, Any]
|
||||||
|
] = tx.clone()
|
||||||
|
|
||||||
|
bstream: trionics.BroadcastReceiver = await nurse.start(
|
||||||
|
relay_to_common_memchan,
|
||||||
|
stream,
|
||||||
|
relay_tx,
|
||||||
|
)
|
||||||
|
broker2bstreams[broker] = bstream
|
||||||
|
|
||||||
|
# Each relay now owns one `tx` clone. Closing the original
|
||||||
|
# lets `rx` reach EOC only after every provider relay exits.
|
||||||
|
await tx.aclose()
|
||||||
try:
|
try:
|
||||||
yield rx
|
yield rx
|
||||||
finally:
|
finally:
|
||||||
|
|
@ -757,10 +810,23 @@ class Feed(Struct):
|
||||||
_max_sample_rate: int = 1
|
_max_sample_rate: int = 1
|
||||||
|
|
||||||
async def pause(self) -> None:
|
async def pause(self) -> None:
|
||||||
|
'''
|
||||||
|
Pause the shared remote provider context for this cached feed.
|
||||||
|
|
||||||
|
This controls provider publication for every local cache user;
|
||||||
|
it does not pause only this caller's receive cursor.
|
||||||
|
|
||||||
|
'''
|
||||||
|
stream: tractor.MsgStream|trionics.BroadcastReceiver
|
||||||
for stream in set(self.streams.values()):
|
for stream in set(self.streams.values()):
|
||||||
await stream.send('pause')
|
await stream.send('pause')
|
||||||
|
|
||||||
async def resume(self) -> None:
|
async def resume(self) -> None:
|
||||||
|
'''
|
||||||
|
Resume the shared remote provider context for this cached feed.
|
||||||
|
|
||||||
|
'''
|
||||||
|
stream: tractor.MsgStream|trionics.BroadcastReceiver
|
||||||
for stream in set(self.streams.values()):
|
for stream in set(self.streams.values()):
|
||||||
await stream.send('resume')
|
await stream.send('resume')
|
||||||
|
|
||||||
|
|
@ -820,7 +886,17 @@ async def maybe_open_feed(
|
||||||
in a tractor broadcast receiver.
|
in a tractor broadcast receiver.
|
||||||
|
|
||||||
'''
|
'''
|
||||||
fqme = fqmes[0]
|
cache_key: tuple[
|
||||||
|
tuple[str, ...],
|
||||||
|
float|None,
|
||||||
|
bool,
|
||||||
|
bool,
|
||||||
|
] = (
|
||||||
|
tuple(sorted(fqmes)),
|
||||||
|
kwargs.get('tick_throttle'),
|
||||||
|
kwargs.get('allow_overruns', True),
|
||||||
|
kwargs.get('start_stream', True),
|
||||||
|
)
|
||||||
|
|
||||||
async with trionics.maybe_open_context(
|
async with trionics.maybe_open_context(
|
||||||
acm_func=open_feed,
|
acm_func=open_feed,
|
||||||
|
|
@ -833,27 +909,73 @@ async def maybe_open_feed(
|
||||||
'allow_overruns': kwargs.get('allow_overruns', True),
|
'allow_overruns': kwargs.get('allow_overruns', True),
|
||||||
'start_stream': kwargs.get('start_stream', True),
|
'start_stream': kwargs.get('start_stream', True),
|
||||||
},
|
},
|
||||||
key=fqme,
|
key=cache_key,
|
||||||
|
|
||||||
) as (cache_hit, feed):
|
) as (cache_hit, feed):
|
||||||
|
feed: Feed
|
||||||
|
|
||||||
if cache_hit:
|
if (
|
||||||
log.info(f'Using cached feed for {fqme}')
|
cache_hit
|
||||||
# add a new broadcast subscription for the quote stream
|
and
|
||||||
# if this feed is likely already in use
|
log.at_least_level('info')
|
||||||
|
):
|
||||||
|
log.info(
|
||||||
|
f'Using cached feed for key:\n'
|
||||||
|
f'{cache_key!r}'
|
||||||
|
)
|
||||||
|
|
||||||
|
# The cached `Feed` exclusively owns provider root streams.
|
||||||
|
# Every caller, including the first cache owner, gets a local
|
||||||
|
# child cursor and copied descriptors for its lexical lifetime.
|
||||||
|
broker_names: tuple[str, ...] = tuple(feed.streams)
|
||||||
|
bstreams: tuple[
|
||||||
|
trionics.BroadcastReceiver,
|
||||||
|
...,
|
||||||
|
]
|
||||||
async with trionics.gather_contexts(
|
async with trionics.gather_contexts(
|
||||||
mngrs=[stream.subscribe() for stream in feed.streams.values()]
|
mngrs=[
|
||||||
|
feed.streams[broker].subscribe()
|
||||||
|
for broker in broker_names
|
||||||
|
]
|
||||||
) as bstreams:
|
) as bstreams:
|
||||||
for bstream, flume in zip(bstreams, feed.flumes.values()):
|
broker2bstreams: dict[
|
||||||
|
str,
|
||||||
|
trionics.BroadcastReceiver,
|
||||||
|
] = dict(zip(
|
||||||
|
broker_names,
|
||||||
|
bstreams,
|
||||||
|
strict=True,
|
||||||
|
))
|
||||||
|
flumes: dict[str, Flume] = {}
|
||||||
|
fqme: str
|
||||||
|
flume: Flume
|
||||||
|
for fqme, flume in feed.flumes.items():
|
||||||
|
broker: str = flume.mkt.broker
|
||||||
|
|
||||||
|
bstream: trionics.BroadcastReceiver = (
|
||||||
|
broker2bstreams[broker]
|
||||||
|
)
|
||||||
|
|
||||||
# XXX: TODO: horrible hackery that needs fixing..
|
# XXX: TODO: horrible hackery that needs fixing..
|
||||||
# i guess we have to create context proxies?
|
# i guess we have to create context proxies?
|
||||||
bstream._ctx = flume.stream._ctx
|
root_stream: tractor.MsgStream = feed.streams[broker]
|
||||||
flume.stream = bstream
|
|
||||||
|
|
||||||
yield feed
|
ctx: tractor.Context = root_stream.ctx
|
||||||
else:
|
bstream._ctx = ctx
|
||||||
yield feed
|
|
||||||
|
local_flume: Flume = copy(flume)
|
||||||
|
local_flume.stream = bstream
|
||||||
|
flumes[fqme] = local_flume
|
||||||
|
|
||||||
|
# `feed` is the cache-owned resource and must retain only
|
||||||
|
# provider root streams. A shallow `Feed` copy gives this
|
||||||
|
# caller its own descriptor mappings while intentionally
|
||||||
|
# sharing immutable market metadata, SHM handles and
|
||||||
|
# provider controls for the cache resource's lifetime.
|
||||||
|
local_feed: Feed = copy(feed)
|
||||||
|
local_feed.flumes = flumes
|
||||||
|
local_feed.streams = broker2bstreams
|
||||||
|
yield local_feed
|
||||||
|
|
||||||
|
|
||||||
@acm
|
@acm
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,112 @@
|
||||||
|
'''
|
||||||
|
Deribit cached feed broadcast ownership regressions.
|
||||||
|
|
||||||
|
'''
|
||||||
|
from contextlib import asynccontextmanager as acm
|
||||||
|
|
||||||
|
import tractor
|
||||||
|
import trio
|
||||||
|
|
||||||
|
from piker.brokers.deribit import api
|
||||||
|
|
||||||
|
|
||||||
|
def test_price_feed_owns_child_subscription(
|
||||||
|
monkeypatch,
|
||||||
|
) -> None:
|
||||||
|
'''
|
||||||
|
Every cached price-feed caller must own one lexical child receiver.
|
||||||
|
|
||||||
|
`maybe_open_price_feed()` previously yielded the raw
|
||||||
|
`LinkedTaskChannel` to the first cache owner, then constructed a new
|
||||||
|
broadcaster root around that same channel for every cache hit.
|
||||||
|
Concurrent roots competed for source receives, so quotes could be
|
||||||
|
split across callers and exiting one wrapper did not express child
|
||||||
|
ownership on the retained channel broadcaster.
|
||||||
|
|
||||||
|
Return one real `LinkedTaskChannel`-shaped fake first as a miss and
|
||||||
|
then as a concurrent hit. Nest both caller contexts and prove each
|
||||||
|
gets a distinct child from the channel's single retained root and
|
||||||
|
receives the same quote. Close the hit child, then prove the first
|
||||||
|
owner continues receiving before its own lexical child closes.
|
||||||
|
|
||||||
|
'''
|
||||||
|
class FakeLinkedTaskChannel:
|
||||||
|
'''
|
||||||
|
Retain one real broadcaster and expose its public subscription.
|
||||||
|
|
||||||
|
'''
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self._tx: trio.MemorySendChannel
|
||||||
|
rx: trio.MemoryReceiveChannel
|
||||||
|
self._tx, rx = trio.open_memory_channel(8)
|
||||||
|
self._broadcaster: (
|
||||||
|
tractor.trionics.BroadcastReceiver
|
||||||
|
) = tractor.trionics.broadcast_receiver(
|
||||||
|
rx,
|
||||||
|
8,
|
||||||
|
)
|
||||||
|
|
||||||
|
@acm
|
||||||
|
async def subscribe(self):
|
||||||
|
'''
|
||||||
|
Yield one caller-owned child from the retained root.
|
||||||
|
|
||||||
|
'''
|
||||||
|
child: tractor.trionics.BroadcastReceiver
|
||||||
|
async with self._broadcaster.subscribe() as child:
|
||||||
|
yield child
|
||||||
|
|
||||||
|
async def push(self, msg: dict) -> None:
|
||||||
|
'''
|
||||||
|
Send one quote through the retained source channel.
|
||||||
|
|
||||||
|
'''
|
||||||
|
await self._tx.send(msg)
|
||||||
|
|
||||||
|
channel: FakeLinkedTaskChannel = FakeLinkedTaskChannel()
|
||||||
|
context_entries: int = 0
|
||||||
|
|
||||||
|
@acm
|
||||||
|
async def maybe_open_context(**kwargs):
|
||||||
|
'''
|
||||||
|
Return the controlled channel with either cache state.
|
||||||
|
|
||||||
|
'''
|
||||||
|
nonlocal context_entries
|
||||||
|
cache_hit: bool = context_entries > 0
|
||||||
|
context_entries += 1
|
||||||
|
yield cache_hit, channel
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
api,
|
||||||
|
'maybe_open_context',
|
||||||
|
maybe_open_context,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def main() -> None:
|
||||||
|
owner: tractor.trionics.BroadcastReceiver
|
||||||
|
async with api.maybe_open_price_feed(
|
||||||
|
'btc-usd',
|
||||||
|
) as owner:
|
||||||
|
hit: tractor.trionics.BroadcastReceiver
|
||||||
|
async with api.maybe_open_price_feed(
|
||||||
|
'btc-usd',
|
||||||
|
) as hit:
|
||||||
|
assert hit is not owner
|
||||||
|
assert len(channel._broadcaster._state.subs) == 3
|
||||||
|
|
||||||
|
await channel.push({'value': 1})
|
||||||
|
assert await owner.receive() == {'value': 1}
|
||||||
|
assert await hit.receive() == {'value': 1}
|
||||||
|
|
||||||
|
assert hit._closed
|
||||||
|
assert not owner._closed
|
||||||
|
assert len(channel._broadcaster._state.subs) == 2
|
||||||
|
|
||||||
|
await channel.push({'value': 2})
|
||||||
|
assert await owner.receive() == {'value': 2}
|
||||||
|
|
||||||
|
assert owner._closed
|
||||||
|
assert len(channel._broadcaster._state.subs) == 1
|
||||||
|
|
||||||
|
trio.run(main)
|
||||||
|
|
@ -0,0 +1,79 @@
|
||||||
|
'''
|
||||||
|
EMS broadcast ownership regressions.
|
||||||
|
|
||||||
|
'''
|
||||||
|
from contextlib import asynccontextmanager as acm
|
||||||
|
|
||||||
|
import trio
|
||||||
|
|
||||||
|
from piker.clearing import _ems
|
||||||
|
|
||||||
|
|
||||||
|
def test_dark_clearing_owns_non_strict_subscription(
|
||||||
|
monkeypatch,
|
||||||
|
) -> None:
|
||||||
|
'''
|
||||||
|
Dark clearing must not mutate another feed consumer's lag policy.
|
||||||
|
|
||||||
|
`clear_dark_triggers()` previously assigned `_raise_on_lag` on
|
||||||
|
whichever root or temporary receiver was stored in `Flume.stream`.
|
||||||
|
That changed private shared state and made behavior depend on feed
|
||||||
|
cache ownership. Supply a stream whose subscription records the
|
||||||
|
public lag-policy argument and yields a distinct child. Replace
|
||||||
|
the hot trigger loop with a checkpointing probe, then prove the
|
||||||
|
non-strict child is entered, passed to the core, and closed only
|
||||||
|
after processing completes while the source remains untouched.
|
||||||
|
|
||||||
|
'''
|
||||||
|
events: list[tuple] = []
|
||||||
|
child = object()
|
||||||
|
|
||||||
|
class QuoteStream:
|
||||||
|
'''
|
||||||
|
Record public subscription ownership without quote traffic.
|
||||||
|
|
||||||
|
'''
|
||||||
|
@acm
|
||||||
|
async def subscribe(self, raise_on_lag: bool = True):
|
||||||
|
'''
|
||||||
|
Yield one dedicated child for the wrapper lifetime.
|
||||||
|
|
||||||
|
'''
|
||||||
|
events.append(('enter', raise_on_lag))
|
||||||
|
try:
|
||||||
|
yield child
|
||||||
|
finally:
|
||||||
|
events.append(('exit', raise_on_lag))
|
||||||
|
|
||||||
|
async def clear_core(*args) -> None:
|
||||||
|
'''
|
||||||
|
Prove the child remains owned while processing runs.
|
||||||
|
|
||||||
|
'''
|
||||||
|
events.append(('core', args[2] is child))
|
||||||
|
await trio.lowlevel.checkpoint()
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
_ems,
|
||||||
|
'_clear_dark_triggers',
|
||||||
|
clear_core,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def main() -> None:
|
||||||
|
stream = QuoteStream()
|
||||||
|
await _ems.clear_dark_triggers(
|
||||||
|
object(),
|
||||||
|
object(),
|
||||||
|
stream,
|
||||||
|
'ib',
|
||||||
|
'nvda.nasdaq.ib',
|
||||||
|
object(),
|
||||||
|
)
|
||||||
|
assert not hasattr(stream, '_raise_on_lag')
|
||||||
|
|
||||||
|
trio.run(main)
|
||||||
|
assert events == [
|
||||||
|
('enter', False),
|
||||||
|
('core', True),
|
||||||
|
('exit', False),
|
||||||
|
]
|
||||||
|
|
@ -0,0 +1,347 @@
|
||||||
|
'''
|
||||||
|
Cached feed broadcast ownership regressions.
|
||||||
|
|
||||||
|
'''
|
||||||
|
from contextlib import asynccontextmanager as acm
|
||||||
|
from types import SimpleNamespace
|
||||||
|
|
||||||
|
import tractor
|
||||||
|
import trio
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from piker.data import feed as feed_mod
|
||||||
|
from piker.data.flows import Flume
|
||||||
|
|
||||||
|
|
||||||
|
class FakeMsgStream:
|
||||||
|
'''
|
||||||
|
Duplex test stream backed by a real Tractor broadcaster.
|
||||||
|
|
||||||
|
'''
|
||||||
|
def __init__(self, broker: str) -> None:
|
||||||
|
self.broker: str = broker
|
||||||
|
self._ctx: object = object()
|
||||||
|
self._tx, rx = trio.open_memory_channel(8)
|
||||||
|
self._broadcaster = tractor.trionics.broadcast_receiver(
|
||||||
|
rx,
|
||||||
|
8,
|
||||||
|
)
|
||||||
|
self.children: list[
|
||||||
|
tractor.trionics.BroadcastReceiver,
|
||||||
|
] = []
|
||||||
|
self.controls: list[str] = []
|
||||||
|
|
||||||
|
@property
|
||||||
|
def ctx(self) -> object:
|
||||||
|
'''
|
||||||
|
Expose the public context ref provided by `MsgStream.ctx`.
|
||||||
|
|
||||||
|
'''
|
||||||
|
return self._ctx
|
||||||
|
|
||||||
|
@acm
|
||||||
|
async def subscribe(self):
|
||||||
|
'''
|
||||||
|
Yield a real child receiver with duplex send delegation.
|
||||||
|
|
||||||
|
'''
|
||||||
|
child: tractor.trionics.BroadcastReceiver
|
||||||
|
async with self._broadcaster.subscribe() as child:
|
||||||
|
child.send = self.send
|
||||||
|
self.children.append(child)
|
||||||
|
yield child
|
||||||
|
|
||||||
|
async def send(self, msg: str) -> None:
|
||||||
|
'''
|
||||||
|
Record a pause or resume control sent through a child.
|
||||||
|
|
||||||
|
'''
|
||||||
|
self.controls.append(msg)
|
||||||
|
|
||||||
|
async def push(self, msg: dict) -> None:
|
||||||
|
'''
|
||||||
|
Send one provider quote into the broadcast source.
|
||||||
|
|
||||||
|
'''
|
||||||
|
await self._tx.send(msg)
|
||||||
|
|
||||||
|
async def close_source(self) -> None:
|
||||||
|
'''
|
||||||
|
End this provider without closing sibling relay channels.
|
||||||
|
|
||||||
|
'''
|
||||||
|
await self._tx.aclose()
|
||||||
|
|
||||||
|
|
||||||
|
def make_feed() -> tuple[
|
||||||
|
feed_mod.Feed,
|
||||||
|
dict[str, FakeMsgStream],
|
||||||
|
]:
|
||||||
|
'''
|
||||||
|
Build a real `Feed` and `Flume` graph for two providers.
|
||||||
|
|
||||||
|
'''
|
||||||
|
roots: dict[str, FakeMsgStream] = {
|
||||||
|
broker: FakeMsgStream(broker)
|
||||||
|
for broker in ('ib', 'deribit')
|
||||||
|
}
|
||||||
|
flumes: dict[str, Flume] = {
|
||||||
|
fqme: Flume(
|
||||||
|
mkt=SimpleNamespace(broker=broker),
|
||||||
|
first_quote={},
|
||||||
|
_rt_shm_token=SimpleNamespace(),
|
||||||
|
stream=roots[broker],
|
||||||
|
)
|
||||||
|
for fqme, broker in (
|
||||||
|
('nvda.nasdaq.ib', 'ib'),
|
||||||
|
('aapl.nasdaq.ib', 'ib'),
|
||||||
|
('btc-usd.deribit', 'deribit'),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
return (
|
||||||
|
feed_mod.Feed(
|
||||||
|
mods={
|
||||||
|
broker: SimpleNamespace(name=broker)
|
||||||
|
for broker in roots
|
||||||
|
},
|
||||||
|
portals={},
|
||||||
|
flumes=flumes,
|
||||||
|
streams=roots,
|
||||||
|
status={},
|
||||||
|
),
|
||||||
|
roots,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize('cache_hit', [False, True])
|
||||||
|
def test_cached_feed_returns_local_streams(
|
||||||
|
monkeypatch,
|
||||||
|
cache_hit: bool,
|
||||||
|
) -> None:
|
||||||
|
'''
|
||||||
|
A cached feed caller must not consume or replace shared streams.
|
||||||
|
|
||||||
|
`maybe_open_feed()` previously paired provider subscriptions with
|
||||||
|
`Feed.flumes.values()` using `zip()` and installed each child on the
|
||||||
|
cached `Flume`. One caller therefore exposed its lexical receiver
|
||||||
|
to every holder of the shared `Feed`; after exit those descriptors
|
||||||
|
pointed at closed children. The positional pairing also dropped a
|
||||||
|
market whenever one provider supplied multiple flumes.
|
||||||
|
|
||||||
|
Model two real provider broadcasters with two IB markets and one
|
||||||
|
Deribit market. Enter both first-owner and cache-hit paths, then
|
||||||
|
prove the yielded `Feed`, its `Feed.streams` mapping and all
|
||||||
|
`Flume`s are caller-local views. Exercise provider-context-wide
|
||||||
|
pause/resume and `Feed.open_multi_stream()` to show duplex sends
|
||||||
|
still reach the roots while each relay owns another child. Finally
|
||||||
|
prove every lexical child closes and all cached descriptors still
|
||||||
|
point at live roots.
|
||||||
|
|
||||||
|
'''
|
||||||
|
cached_feed: feed_mod.Feed
|
||||||
|
roots: dict[str, FakeMsgStream]
|
||||||
|
cached_feed, roots = make_feed()
|
||||||
|
cached_flumes: dict[str, Flume] = cached_feed.flumes
|
||||||
|
|
||||||
|
@acm
|
||||||
|
async def maybe_open_context(**kwargs):
|
||||||
|
'''
|
||||||
|
Return the controlled shared feed as a cache hit.
|
||||||
|
|
||||||
|
'''
|
||||||
|
yield cache_hit, cached_feed
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
feed_mod.trionics,
|
||||||
|
'maybe_open_context',
|
||||||
|
maybe_open_context,
|
||||||
|
)
|
||||||
|
|
||||||
|
local_feed: feed_mod.Feed|None = None
|
||||||
|
local_streams: dict[
|
||||||
|
str,
|
||||||
|
tractor.trionics.BroadcastReceiver,
|
||||||
|
]|None = None
|
||||||
|
|
||||||
|
async def main() -> None:
|
||||||
|
nonlocal local_feed, local_streams
|
||||||
|
feed: feed_mod.Feed
|
||||||
|
async with feed_mod.maybe_open_feed(
|
||||||
|
list(cached_flumes),
|
||||||
|
) as feed:
|
||||||
|
local_feed = feed
|
||||||
|
local_streams = feed.streams
|
||||||
|
assert feed is not cached_feed
|
||||||
|
assert feed.streams is not roots
|
||||||
|
assert feed.flumes is not cached_flumes
|
||||||
|
|
||||||
|
broker: str
|
||||||
|
stream: tractor.trionics.BroadcastReceiver
|
||||||
|
for broker, stream in feed.streams.items():
|
||||||
|
assert stream is roots[broker].children[0]
|
||||||
|
assert stream._ctx is roots[broker].ctx
|
||||||
|
|
||||||
|
fqme: str
|
||||||
|
flume: Flume
|
||||||
|
for fqme, flume in feed.flumes.items():
|
||||||
|
cached: Flume = cached_flumes[fqme]
|
||||||
|
|
||||||
|
stream: tractor.trionics.BroadcastReceiver = (
|
||||||
|
feed.streams[flume.mkt.broker]
|
||||||
|
)
|
||||||
|
|
||||||
|
assert flume is not cached
|
||||||
|
assert flume.stream is stream
|
||||||
|
assert cached.stream is roots[flume.mkt.broker]
|
||||||
|
|
||||||
|
assert (
|
||||||
|
feed.flumes['nvda.nasdaq.ib'].stream
|
||||||
|
is feed.flumes['aapl.nasdaq.ib'].stream
|
||||||
|
)
|
||||||
|
|
||||||
|
await feed.pause()
|
||||||
|
await feed.resume()
|
||||||
|
assert all(
|
||||||
|
root.controls == ['pause', 'resume']
|
||||||
|
for root in roots.values()
|
||||||
|
)
|
||||||
|
|
||||||
|
stream: trio.MemoryReceiveChannel
|
||||||
|
async with feed.open_multi_stream() as stream:
|
||||||
|
assert all(
|
||||||
|
len(child._state.subs) == 3
|
||||||
|
for child in feed.streams.values()
|
||||||
|
)
|
||||||
|
await roots['ib'].push({'provider': 'ib'})
|
||||||
|
await roots['deribit'].push({
|
||||||
|
'provider': 'deribit',
|
||||||
|
})
|
||||||
|
with trio.fail_after(1):
|
||||||
|
msgs = {
|
||||||
|
(await stream.receive())['provider'],
|
||||||
|
(await stream.receive())['provider'],
|
||||||
|
}
|
||||||
|
assert msgs == {'ib', 'deribit'}
|
||||||
|
|
||||||
|
# One provider's EOC must close only its `tx` clone;
|
||||||
|
# the sibling relay continues publishing to `stream`.
|
||||||
|
await roots['ib'].close_source()
|
||||||
|
with trio.fail_after(1):
|
||||||
|
while (
|
||||||
|
len(feed.streams['ib']._state.subs)
|
||||||
|
!= 2
|
||||||
|
):
|
||||||
|
await trio.lowlevel.checkpoint()
|
||||||
|
|
||||||
|
await roots['deribit'].push({
|
||||||
|
'provider': 'deribit-after-ib-eoc',
|
||||||
|
})
|
||||||
|
msg: dict
|
||||||
|
msg = await stream.receive()
|
||||||
|
|
||||||
|
assert msg['provider'] == 'deribit-after-ib-eoc'
|
||||||
|
|
||||||
|
assert all(
|
||||||
|
len(child._state.subs) == 2
|
||||||
|
for child in feed.streams.values()
|
||||||
|
)
|
||||||
|
|
||||||
|
trio.run(main)
|
||||||
|
|
||||||
|
assert local_feed is not None
|
||||||
|
assert local_streams is not None
|
||||||
|
assert all(stream._closed for stream in local_streams.values())
|
||||||
|
assert all(
|
||||||
|
flume.stream is roots[flume.mkt.broker]
|
||||||
|
for flume in cached_flumes.values()
|
||||||
|
)
|
||||||
|
assert all(
|
||||||
|
len(root._broadcaster._state.subs) == 1
|
||||||
|
for root in roots.values()
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_feed_cache_key_includes_shape_and_policy(
|
||||||
|
monkeypatch,
|
||||||
|
) -> None:
|
||||||
|
'''
|
||||||
|
Feed cache identity must include every resource-defining input.
|
||||||
|
|
||||||
|
The old `fqmes[0]` key aliased `[A]` with `[A, B]` and ignored
|
||||||
|
`tick_throttle`, `allow_overruns` and `start_stream`. Whichever
|
||||||
|
caller entered first silently determined later feed shape and
|
||||||
|
stream behavior, including reuse of a history-only remote feed for
|
||||||
|
a live quote request.
|
||||||
|
|
||||||
|
Capture the key passed to `maybe_open_context()` for reordered and
|
||||||
|
varied requests. Symbol order must normalize to one key, while the
|
||||||
|
full symbol set and each stream policy must produce distinct keys.
|
||||||
|
|
||||||
|
'''
|
||||||
|
feed: feed_mod.Feed = feed_mod.Feed(
|
||||||
|
mods={},
|
||||||
|
portals={},
|
||||||
|
flumes={},
|
||||||
|
streams={},
|
||||||
|
status={},
|
||||||
|
)
|
||||||
|
keys: list[tuple] = []
|
||||||
|
|
||||||
|
@acm
|
||||||
|
async def maybe_open_context(**kwargs):
|
||||||
|
'''
|
||||||
|
Capture cache identity without opening a remote feed.
|
||||||
|
|
||||||
|
'''
|
||||||
|
keys.append(kwargs['key'])
|
||||||
|
yield False, feed
|
||||||
|
|
||||||
|
@acm
|
||||||
|
async def gather_contexts(mngrs):
|
||||||
|
'''
|
||||||
|
Permit the metadata-only empty feed used by this key probe.
|
||||||
|
|
||||||
|
'''
|
||||||
|
assert not mngrs
|
||||||
|
yield ()
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
feed_mod.trionics,
|
||||||
|
'maybe_open_context',
|
||||||
|
maybe_open_context,
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
feed_mod.trionics,
|
||||||
|
'gather_contexts',
|
||||||
|
gather_contexts,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def open_once(
|
||||||
|
fqmes: list[str],
|
||||||
|
**kwargs,
|
||||||
|
) -> None:
|
||||||
|
async with feed_mod.maybe_open_feed(
|
||||||
|
fqmes,
|
||||||
|
**kwargs,
|
||||||
|
):
|
||||||
|
pass
|
||||||
|
|
||||||
|
async def main() -> None:
|
||||||
|
await open_once(['a.ib', 'b.ib'])
|
||||||
|
await open_once(['b.ib', 'a.ib'])
|
||||||
|
await open_once(['a.ib'])
|
||||||
|
await open_once(['a.ib', 'b.ib'], tick_throttle=10)
|
||||||
|
await open_once(
|
||||||
|
['a.ib', 'b.ib'],
|
||||||
|
allow_overruns=False,
|
||||||
|
)
|
||||||
|
await open_once(
|
||||||
|
['a.ib', 'b.ib'],
|
||||||
|
start_stream=False,
|
||||||
|
)
|
||||||
|
|
||||||
|
trio.run(main)
|
||||||
|
|
||||||
|
assert keys[0] == keys[1]
|
||||||
|
assert len(set(keys[2:])) == 4
|
||||||
|
assert keys[0] not in keys[2:]
|
||||||
|
|
@ -0,0 +1,120 @@
|
||||||
|
'''
|
||||||
|
IB cached quote-channel ownership regressions.
|
||||||
|
|
||||||
|
'''
|
||||||
|
from contextlib import asynccontextmanager as acm
|
||||||
|
|
||||||
|
import tractor
|
||||||
|
import trio
|
||||||
|
|
||||||
|
from piker.brokers.ib import feed
|
||||||
|
|
||||||
|
|
||||||
|
def test_quote_stream_caches_channel_and_owns_child(
|
||||||
|
monkeypatch,
|
||||||
|
) -> None:
|
||||||
|
'''
|
||||||
|
IB quote callers must share one cached channel through child cursors.
|
||||||
|
|
||||||
|
`open_aio_quote_stream()` previously stored the first caller's raw
|
||||||
|
`LinkedTaskChannel` in `_quote_streams`, yielded that root directly,
|
||||||
|
and wrapped the same source in a new broadcaster for each later
|
||||||
|
caller. The first caller therefore owned source lifetime while
|
||||||
|
competing roots split quote receives between concurrent users.
|
||||||
|
|
||||||
|
Stub `maybe_open_context()` with one retained channel and enter a
|
||||||
|
miss owner plus concurrent hit. Prove both receive the same quote
|
||||||
|
through distinct children of one broadcaster, the hit closes
|
||||||
|
independently, and the owner keeps receiving until its own exit.
|
||||||
|
Also verify both calls request the same symbol cache key.
|
||||||
|
|
||||||
|
'''
|
||||||
|
class FakeLinkedTaskChannel:
|
||||||
|
'''
|
||||||
|
Retain one real broadcaster and expose child subscriptions.
|
||||||
|
|
||||||
|
'''
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self._tx: trio.MemorySendChannel
|
||||||
|
rx: trio.MemoryReceiveChannel
|
||||||
|
self._tx, rx = trio.open_memory_channel(8)
|
||||||
|
self._broadcaster: (
|
||||||
|
tractor.trionics.BroadcastReceiver
|
||||||
|
) = tractor.trionics.broadcast_receiver(
|
||||||
|
rx,
|
||||||
|
8,
|
||||||
|
)
|
||||||
|
|
||||||
|
@acm
|
||||||
|
async def subscribe(self):
|
||||||
|
'''
|
||||||
|
Yield one caller-owned child from the retained root.
|
||||||
|
|
||||||
|
'''
|
||||||
|
child: tractor.trionics.BroadcastReceiver
|
||||||
|
async with self._broadcaster.subscribe() as child:
|
||||||
|
yield child
|
||||||
|
|
||||||
|
async def push(self, msg: dict) -> None:
|
||||||
|
'''
|
||||||
|
Send one ticker through the retained source channel.
|
||||||
|
|
||||||
|
'''
|
||||||
|
await self._tx.send(msg)
|
||||||
|
|
||||||
|
chan: FakeLinkedTaskChannel = FakeLinkedTaskChannel()
|
||||||
|
context_entries: int = 0
|
||||||
|
keys: list[str] = []
|
||||||
|
acm_funcs: list = []
|
||||||
|
|
||||||
|
@acm
|
||||||
|
async def maybe_open_context(**kwargs):
|
||||||
|
'''
|
||||||
|
Return one retained channel as a miss followed by a hit.
|
||||||
|
|
||||||
|
'''
|
||||||
|
nonlocal context_entries
|
||||||
|
keys.append(kwargs['key'])
|
||||||
|
acm_funcs.append(kwargs['acm_func'])
|
||||||
|
cache_hit: bool = context_entries > 0
|
||||||
|
context_entries += 1
|
||||||
|
yield cache_hit, chan
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
feed.tractor.trionics,
|
||||||
|
'maybe_open_context',
|
||||||
|
maybe_open_context,
|
||||||
|
)
|
||||||
|
|
||||||
|
async def main() -> None:
|
||||||
|
owner: tractor.trionics.BroadcastReceiver
|
||||||
|
async with feed.open_aio_quote_stream(
|
||||||
|
'NVDA',
|
||||||
|
) as owner:
|
||||||
|
hit: tractor.trionics.BroadcastReceiver
|
||||||
|
async with feed.open_aio_quote_stream(
|
||||||
|
'NVDA',
|
||||||
|
) as hit:
|
||||||
|
assert hit is not owner
|
||||||
|
assert len(chan._broadcaster._state.subs) == 3
|
||||||
|
|
||||||
|
await chan.push({'value': 1})
|
||||||
|
assert await owner.receive() == {'value': 1}
|
||||||
|
assert await hit.receive() == {'value': 1}
|
||||||
|
|
||||||
|
assert hit._closed
|
||||||
|
assert not owner._closed
|
||||||
|
assert len(chan._broadcaster._state.subs) == 2
|
||||||
|
|
||||||
|
await chan.push({'value': 2})
|
||||||
|
assert await owner.receive() == {'value': 2}
|
||||||
|
|
||||||
|
assert owner._closed
|
||||||
|
assert len(chan._broadcaster._state.subs) == 1
|
||||||
|
assert keys == ['NVDA', 'NVDA']
|
||||||
|
assert acm_funcs == [
|
||||||
|
feed._open_aio_quote_channel,
|
||||||
|
feed._open_aio_quote_channel,
|
||||||
|
]
|
||||||
|
|
||||||
|
trio.run(main)
|
||||||
Loading…
Reference in New Issue