Repository navigation
Zero-allocation connector boundary (design 054) - #296
Merged
Merged
Conversation
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ign 054 §4.3) Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Move the poll loop into ReadyRoutes::poll_ready, so the orderings that keep a wake-up from being lost live in one tested function instead of being rules every caller of take_next/served must keep. The caller passes a closure that polls one route with the context it is handed and reports Pending, Staged, Skipped or Closed. poll_ready registers the task before reading any bit, polls each route with its own waker, polls a route again after a skipped value or lag, and calls served only right before returning Ready. Getting any of these wrong stalls a route for good, since a reader that returned a value keeps no waker. - register, waker, pass, take_next, served, close and is_done are now private; any_ready is test-only. - Index directly where ids come from 0..len instead of turning an id bug into a silently dropped wake. - Document why Pass bounds a call (Tokio readers wake themselves once the budget is spent) and why set/clear must be single RMWs. - RouteId moves to outbound/mod.rs. - expect(dead_code) instead of allow, so the attribute must go once OutboundRoutes uses the module. Tests drive poll_ready through a Tokio-like reader model. New: skipped values, a spent budget, a producer racing the poll, a randomized liveness test (300 seeds: writes, skips, closes, spent budgets, racing producers; every value written must go out), and a stress test for neighbouring bits in one word. concurrent_wakes_lose_nothing now fails when a park times out with work outstanding instead of rescanning. Each of eight seeded bugs (no served, moving on after a skip, register after the scan, clear after the poll, polling with the task's waker, no task wake, load-then-store clear, no pass bound) fails at least one test; before, the concurrency test passed with the task wake deleted. Design 054: §4.2 says poll_ready owns the loop; §7 adds SelectAll and FuturesUnordered as alternative 13. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SHcxmvGhBYjDxU8wSzNpUQ
… in the ready set Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A route whose values all fail to stage kept poll_ready polling it for as long as its producer kept writing. Only Tokio's coop budget bounded that loop, and the Tokio mailbox reader and every Embassy reader do not use it, so one route's topic overflow could stall every route of the connector and the rest of its loop. - poll_ready polls a skipping route at most SKIP_BUDGET (32) times in a row; then the route's own waker sets its bit and wakes the task, and it is polled again in a later pass. No route is left with a clear bit and no waker. - The transport task must not outrank its producers (module doc, design 054 §4.2). AtomicWaker::register answers a wake in progress by waking the task again, so a transport that preempted that wake livelocks on one core; reproduced with both threads pinned to one CPU under SCHED_FIFO. Design 049's deployment already satisfies the rule. - Module doc: a buffer that keeps a route waker after its reader is gone frees the set, and wakes the finished task, from the producer's context on its next write. - Tests: a route that keeps skipping does not hold the call; a pass that starts at bit 31 of a word takes that route once. An off-by-one in that word's mask passed every existing test. - Design 054: the §4.2 sketch uses the code's names; §4.2 and §4.6 state the skip bound; alternative 13 no longer counts the half-updated queue against FuturesUnordered, since AtomicWaker shares it; alternative 14 covers futures-concurrency's Merge. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KTnKDhY6tvu8bnscus1TZk
ReadyRoutes::new computed each word of the open mask from the number of routes left, with a special case for a partly filled last word. It now sets each route's bit with the same bit() the wakers and the scan use, so the mask is right by construction. The result is identical for every route count from 0 to 1,024. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KTnKDhY6tvu8bnscus1TZk
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…r rules - poll_stage logs every skip and close with the route's default topic. A serializer failure used to name only the value's type, and a closed route was logged at debug without its error; the pumps log a stop at info with the error, and OutboundRoutes now does too. - The topic-overflow rule and the scratch, fallback and invalid-length rules live in write_topic and serialize_outbound, which FusedReader and TypedRoute both call, so the pumps and OutboundRoutes cannot drift before stage 16. - AimDb::outbound_links holds the storage-index debug_assert, and collect_outbound_routes is built on it. - record_index is set on the parsed config instead of round-tripping through the query string. - OutboundRoutes: Send is asserted at compile time; a test pulls with poll_next from a spawned task. - Docs: a link's profiling interval under OutboundRoutes also covers serving the other ready routes; route_factory is set in finish(); the shared consumer factory's panics no longer say "source factory". Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01VJdtLXS2xKnn3hY3huLtjX
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… (design 054 §4.7) (#286)
…sign 054 §5, §8) (#290)
Records #284 (Maximum Packet Size in the embedded CONNECT) as merged. Its content already reached the feature branch through #295, whose squash dropped the merge parent, so #296 conflicts with main. The resolution keeps the feature branch's side; the tree is unchanged. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- Introduced a new SVG file depicting the AimDB architecture. - The diagram illustrates the flow of data through inbound and outbound connectors, detailing the roles of aimdb CLI, MCP server, and AimX server. - Included annotations for various components such as TCP, Unix socket, Serial, and WebSocket connections. - Visual representation of the AIMDB NODE and its record handling mechanisms.
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2 of 4 tasks
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Design 054, the zero-allocation connector boundary (
docs/design/054-zero-alloc-connector-boundary.md). It lands stages 3–16, each already reviewed as its own PR intofeat/054-connector-boundary(#277–#295).Breaking. Stage 17, the CHANGELOGs and user-facing docs, follows after the final review of this PR.
What changes
Connectors no longer get per-route pump tasks from core. Each connector builds two objects in
build()and moves them into its own transport task:InboundDispatch(new(db, scheme, grammar)) —dispatch(topic, payload)deserializes and produces into every matching record. It is synchronous, never blocks, and does not allocate for exact, pattern or known-key topics.OutboundRoutes(new(db, scheme)) — the transport pulls withnext(), or withpoll_stage+take_stagedinside aselect, when it can send.OutboundPayload::Borrowed); owned serializers fall back toOwned.RouteStatscounts sent, rejected, lagged, topic-overflow and serialize-failed per route.TopicWriter/with_topic_writer/with_topic_fn— write the topic into a boundedTopicBufand replaceTopicProvider.Per crate
aimdb-core:InboundDispatch,OutboundRoutes(plusRouteInfo,OutboundMessage,OutboundPayload,RouteStats),TopicWriter/TopicBuf/TopicOverflow,Reader::poll_recv.pump_clientnow pulls fromOutboundRoutesin one task.Source,pump_source,pump_sink,Connector, theSerialized*types,TopicProvider,with_topic_provider,collect_outbound_routes/OutboundRoute.Routerandinbound_router.aimdb-mqtt-connector:bbqueuewrite ring, sized at build (with_write_buffer, default 4096). Inbound PUBLISHes are dispatched from the session loop and outbound values are pulled fromOutboundRoutes.qos/retainfails the build.aimdb-websocket-connector: the server dispatches client writes and pulls broadcasts in one task.aimdb-knx-connector:KnxConnector::new(binder, delay, gateway_url);Channelsand the related types are gone.critical-section-std-implis a deprecated no-op.aimdb-embassy-adapter: theEmbassySink/EmbassySourcebridges are removed;into_box_futureandNetStackstay.aimdb-codegentemplates and docs: updated to the new constructors.Measured
b0_alloc_connector, gated bymake bench-gatein CI:TopicProvider)MQTT round trip (
alloc_round_trip.rs): the embedded backend allocates 0 times per round trip; the native backend stays at or below 11 (rumqttc).Testing
main).Before merging
🤖 Generated with Claude Code