diff --git a/README.md b/README.md index 7851dee..ab0a27d 100644 --- a/README.md +++ b/README.md @@ -18,6 +18,7 @@ Aligned with upstream `packages/core` (4.0 naming: `Fiber`, formerly `EffectScop - **Plugin system** — `ctx.plugin` applies a plugin *as a revertible effect* on the parent fiber, so the whole plugin tree is one nested effect tree; loading is deferred one tick, load/unload transitions are serialized per fiber (inertia lock) - **Reactive coeffects** — `ctx.provide` / `ctx.inject`: consumers load when their dependencies are satisfied, reload when a provider is swapped, and are torn down *before* their provider finishes unloading - **Events** — `ctx.on` / `once` / `emit` / `bail` / `waterfall` (sync) and `ctx.parallel` / `serial` (async); listeners are effects, removed automatically on fiber disposal +- **Isolation & intercept** — `ctx.isolate` puts a service name in its own realm (provides/injects no longer cross the boundary; share a realm by passing the same label); `ctx.intercept` / Hash inject configs carry per-caller config, resolved with `ctx.resolve_config` ```ruby ctx = Cordis::Context.new @@ -44,7 +45,7 @@ end - [x] Plugin/fiber lifecycle (epoch + inertia state machine, via `async`) - [x] Reactive coeffects (`ctx.provide` / `ctx.inject`) - [x] Event system (`ctx.on`, waterfall included) -- [ ] Isolation & intercept (`ctx.isolate` / `ctx.intercept`) +- [x] Isolation & intercept (`ctx.isolate` / `ctx.intercept`) - [ ] `Cordis::Service` base class - [ ] Loader / hot-reload reconciliation — maybe, later @@ -72,6 +73,7 @@ cordis-rb 用 Ruby 重新實作 [Cordis](https://github.com/cordiverse/cordis)(T - **Plugin 系統** — `ctx.plugin` 把 plugin *當成一個 revertible effect* 掛在 parent fiber 上,整棵 plugin tree 就是巢狀 effect tree;載入延後一個 tick,每個 fiber 的 load/unload 由 inertia lock 序列化 - **Reactive coeffect** — `ctx.provide` / `ctx.inject`:依賴滿足才載入、provider 被換掉就 reload、provider 卸載前依賴者先 teardown - **事件系統** — `ctx.on` / `once` / `emit` / `bail` / `waterfall`(同步)與 `ctx.parallel` / `serial`(非同步);listener 就是 effect,fiber 卸載時自動移除 +- **Isolation 與 intercept** — `ctx.isolate` 讓某個 service name 進入獨立 realm(provide/inject 不再跨界;傳同一個 label 可共用 realm);`ctx.intercept` 與 Hash 形式的 inject config 攜帶 per-caller 設定,用 `ctx.resolve_config` 解析 ```ruby ctx = Cordis::Context.new @@ -98,7 +100,7 @@ end - [x] Plugin/fiber 生命週期(epoch + inertia 狀態機,基於 `async`) - [x] Reactive coeffect(`ctx.provide` / `ctx.inject`) - [x] 事件系統(`ctx.on`,含 waterfall) -- [ ] Isolation 與 intercept(`ctx.isolate` / `ctx.intercept`) +- [x] Isolation 與 intercept(`ctx.isolate` / `ctx.intercept`) - [ ] `Cordis::Service` base class - [ ] Loader / hot-reload reconciliation — 骨架穩了再說 @@ -126,6 +128,7 @@ cordis-rb は、[Cordis](https://github.com/cordiverse/cordis)(TypeScript 製の - **Plugin システム** — `ctx.plugin` は plugin を *revertible effect として* parent fiber に掛けるため、plugin ツリー全体が入れ子の effect ツリーになる;ロードは 1 tick 遅延、fiber ごとの load/unload は inertia lock で直列化 - **Reactive coeffect** — `ctx.provide` / `ctx.inject`:依存が満たされたらロード、provider が入れ替われば reload、provider のアンロード前に依存側が先に teardown - **イベントシステム** — `ctx.on` / `once` / `emit` / `bail` / `waterfall`(同期)と `ctx.parallel` / `serial`(非同期);listener は effect であり、fiber の破棄時に自動で外れる +- **Isolation と intercept** — `ctx.isolate` は service name を独立した realm に隔離(provide/inject は境界を越えない;同じ label を渡せば realm を共有);`ctx.intercept` と Hash 形式の inject config は per-caller 設定を運び、`ctx.resolve_config` で解決 ```ruby ctx = Cordis::Context.new @@ -152,7 +155,7 @@ end - [x] Plugin/fiber ライフサイクル(epoch + inertia ステートマシン、`async` ベース) - [x] Reactive coeffect(`ctx.provide` / `ctx.inject`) - [x] イベントシステム(`ctx.on`、waterfall 含む) -- [ ] Isolation と intercept(`ctx.isolate` / `ctx.intercept`) +- [x] Isolation と intercept(`ctx.isolate` / `ctx.intercept`) - [ ] `Cordis::Service` base class - [ ] Loader / hot-reload reconciliation — 骨格が安定してから diff --git a/agents.md b/agents.md index 3e850ee..2b21818 100644 --- a/agents.md +++ b/agents.md @@ -27,6 +27,7 @@ notebooklm ask "..." --notebook b4683395-8795-47c8-9d7a-76f058954f06 - **載入延後一個 tick**(對應上游 `await Promise.resolve()`):async 的 child task 是 eager,所以 transition task 開頭 `Async::Task.current.yield`。 - **inertia lock**:每個 fiber 一個 transition task,loop 到 current epoch == target epoch 為止;transition 中 target 變更只記意圖。`restart` 必須先 `await` 卸載完再 `refresh`,否則 target 設回原 epoch 會把 unload 意圖合併掉(已踩過)。 - **與上游的刻意差異**(`# ponytail:` 註記):unload 是嚴格 LIFO「循序」撤除(上游是 LIFO 啟動 + 併發完成);plugin 回傳值只在 callable 時收為 disposer(上游 `_execute` 會對非法值 TypeError)。 +- **isolate/intercept(round 3)**:isolate key 預設就是 name 本身(不學上游 `root[isolate][name] ??= Symbol` 的 lazy 建 key),isolated ctx 才在自己的 `@isolate` map 放 token(`Object.new` 或使用者傳的 label);map 是 copy-on-write merge,不用 prototype chain。service store 改以 key 為鍵、`Fiber#provided` 改 `{name => key}`、notify 收 `{name => key}` pairs 並過濾 `fib.ctx.isolate_key(n) == key`。intercept 同樣 copy-on-write,鏈在 `ctx.intercept` 時就地 merge 攤平(上游是 resolveConfig 走 prototype chain 收集再 assign),讀取入口是 `ctx.resolve_config(name, base)`;Hash 形式的 inject config 會轉成 plugin ctx 上的 intercept(對齊 fiber.ts:139)。**比上游嚴格的一處**:service walk 在終點 fiber 對 provided 服務檢查 realm key(上游 name-keyed store 會讓 isolated plugin 走到 root 時漏看到 default realm 的 service)。 - async 模式參考 `/Users/ryudo/RubyPrjs/lens-ruby-async`:完成訊號用 `Async::Queue` 不用 `Async::Condition`(signal 先於 wait 會遺失)、spec 用 `task.with_timeout` 包可能吊死的流程(spec_helper 的 `with_reactor`)、teardown 逐步各自 rescue。 ## 範圍紀律 diff --git a/examples/demo.rb b/examples/demo.rb index 810ee12..2269835 100755 --- a/examples/demo.rb +++ b/examples/demo.rb @@ -68,6 +68,24 @@ def banner(text) db_fiber.dispose puts " web state: #{web_fiber.state} (back to pending, waiting for the next provider)" + banner 'isolate(:db): two tenants, same plugin code, separate realms' + alice = ctx.isolate(:db) + bob = ctx.isolate(:db) + alice_web = alice.plugin(web) + bob_web = bob.plugin(web) + alice.plugin(database, { name: 'alice-db' }).await + tick + puts " alice web: #{alice_web.state}, bob web: #{bob_web.state} (bob's realm has no :db yet)" + bob_db = bob.plugin(database, { name: 'bob-db' }) + bob_db.await + tick + ctx.emit('request', '/dashboard') # both tenants answer, each with its own db + + banner 'dispose(bob db): only bob\'s realm tears down' + bob_db.dispose + puts " bob web: #{bob_web.state}, alice web: #{alice_web.state}" + ctx.emit('request', '/dashboard') + banner 'root dispose: the whole tree unwinds in LIFO order' ctx.fiber.dispose puts " registry size: #{ctx.registry.size}" diff --git a/lib/cordis/context.rb b/lib/cordis/context.rb index 9483855..c540ea0 100644 --- a/lib/cordis/context.rb +++ b/lib/cordis/context.rb @@ -13,7 +13,9 @@ class Context def initialize @root = self - @services = {} # name => Impl (no isolate layer, keyed directly by name) + @services = {} # isolate key => Impl (default key is the name itself) + @isolate = {} # name => isolation label; absent = the shared default realm + @intercept = {} # name => per-caller config (copy-on-write, chain pre-merged) @fiber = Fiber.new(self) # root fiber: always active, dispose = restart @registry = Registry.new(self) @events = Events.new(self) @@ -26,6 +28,38 @@ def extend(fiber: @fiber) ctx end + # A ctx view where service `name` lives in its own realm: provides/injects through + # this view no longer see (or are seen by) the default realm. Pass the same label + # to two isolate calls to share one realm between them. + def isolate(name, label = nil) + ctx = extend + ctx.instance_variable_set(:@isolate, @isolate.merge(name.to_sym => label || Object.new)) + ctx + end + + # A ctx view carrying per-caller config for service `name` (upstream ctx.intercept). + # Nested intercepts merge, inner overriding outer (Hash configs only; anything else replaces). + def intercept(name, config) + name = name.to_sym + old = @intercept[name] + merged = old.is_a?(Hash) && config.is_a?(Hash) ? old.merge(config) : config + ctx = extend + ctx.instance_variable_set(:@intercept, @intercept.merge(name => merged)) + ctx + end + + # Effective intercept config for `name` as seen from this ctx, merged over `base` + # (upstream Service[resolveConfig], minus the Config schema merge). + def resolve_config(name, base = nil) + config = @intercept[name.to_sym] + return base if config.nil? + + base.is_a?(Hash) && config.is_a?(Hash) ? base.merge(config) : config + end + + # The store key for `name` in this ctx's realm. + def isolate_key(name) = @isolate[name] || name + # -- effect / plugin -- def effect(label = nil, &) = @fiber.effect(label, &) @@ -53,18 +87,19 @@ def waterfall(name, *, &) = @root.events.waterfall(name, *, &) # (so dependents' disposers can still reach the service during their own teardown). def provide(name, value = nil) name = name.to_sym + key = isolate_key(name) owner = @fiber @fiber.effect("ctx.provide(#{name.inspect})") do - raise ServiceError, "service #{name.inspect} has already been registered" if @root.services.key?(name) + raise ServiceError, "service #{name.inspect} has already been registered" if @root.services.key?(key) impl = Impl.new(name, owner, value) - @root.services[name] = impl + @root.services[key] = impl owner.store[name] = impl # immediately visible to the provider itself - owner.provided << name - @root.notify([name]) if owner.state == :active + owner.provided[name] = key + @root.notify({ name => key }) if owner.state == :active lambda do - @root.services.delete(name) - dependents = @root.notify([name]) + @root.services.delete(key) + dependents = @root.notify({ name => key }) dependents.each do |dep| dep.await rescue StandardError @@ -79,7 +114,7 @@ def provide(name, value = nil) # Strict: the provider fiber must be ACTIVE to be visible (a service whose provider # is still loading is invisible to dependents). def get(name) - impl = @root.services[name.to_sym] + impl = @root.services[isolate_key(name.to_sym)] return nil unless impl && impl.fiber.state == :active impl.value @@ -87,7 +122,7 @@ def get(name) def set(name, value) name = name.to_sym - impl = @root.services[name] + impl = @root.services[isolate_key(name)] raise ServiceError, %(cannot set property "#{name}" without provide) unless impl raise ServiceError, %(cannot set property "#{name}" in multiple fibers) unless impl.fiber.equal?(@fiber) @@ -95,20 +130,22 @@ def set(name, value) end # Dependency-graph update (fully synchronous): linear scan over all fibers, - # re-checking satisfaction and recomputing epochs. Returns the touched fibers - # (provide's teardown uses the list to wait for dependents). - def notify(names) + # re-checking satisfaction and recomputing epochs. Only fibers whose ctx resolves + # the name to the same isolate key are touched. Takes { name => key } pairs and + # returns the touched fibers (provide's teardown uses the list to wait for dependents). + def notify(pairs) touched = [] @registry.runtimes.each do |runtime| runtime.fibers.each do |fib| - next unless names.any? { |n| fib.inject.key?(n) } + hits = pairs.select { |n, key| fib.inject.key?(n) && fib.ctx.isolate_key(n) == key } + next if hits.empty? - names.each { |n| fib.check_impl(n) if fib.inject.key?(n) } + hits.each_key { |n| fib.check_impl(n) } fib.refresh touched << fib end end - names.each { |n| @events.emit('internal/service', n) } + pairs.each_key { |n| @events.emit('internal/service', n) } touched end @@ -121,20 +158,26 @@ def method_missing(name, *args, &block) end def respond_to_missing?(name, include_private = false) - @root.services.key?(name) || super + @root.services.key?(isolate_key(name)) || super end private def resolve_service(name) + key = isolate_key(name) fib = @fiber loop do if fib.store impl = fib.store[name] - return impl.value if impl + # a service this fiber *provides* is only visible from the same realm + # (stricter than upstream, whose walk is name-keyed at the terminal fiber) + return impl.value if impl && (!fib.provided.key?(name) || fib.provided[name] == key) end raise ServiceError, %(cannot get required service "#{name}" in inactive context) if fib.inject.key?(name) - raise ServiceError, %(cannot get property "#{name}" without inject) if fib.root? + # stop at root, or when walking up would cross an isolation boundary + if fib.root? || fib.parent_ctx.isolate_key(name) != key + raise ServiceError, %(cannot get property "#{name}" without inject) + end fib = fib.parent_fiber end diff --git a/lib/cordis/fiber.rb b/lib/cordis/fiber.rb index 5398338..d588529 100644 --- a/lib/cordis/fiber.rb +++ b/lib/cordis/fiber.rb @@ -12,14 +12,15 @@ class Fiber Entry = Struct.new(:disposers, :label) - attr_reader :uid, :ctx, :runtime, :inject, :store, :config, :error, :provided, :parent_fiber, :state + attr_reader :uid, :ctx, :runtime, :inject, :store, :config, :error, :provided, :parent_fiber, :parent_ctx, + :state def initialize(parent_ctx, config: nil, inject: {}, runtime: nil) @config = config @inject = inject @runtime = runtime @disposables = [] - @provided = [] + @provided = {} # name => isolate key (for boundary-crossing notify on state change) @internal_store = {} # currently satisfied dependencies (name => Impl), continuously updated @store = nil # snapshot taken while loading; nil means unloaded @current_epoch = INACTIVE @@ -40,7 +41,11 @@ def initialize(parent_ctx, config: nil, inject: {}, runtime: nil) @uid = parent_ctx.root.registry.next_uid @parent_fiber = parent_ctx.fiber + @parent_ctx = parent_ctx @ctx = parent_ctx.extend(fiber: self) + # inject with a Hash config doubles as an intercept entry on the plugin's ctx + # (mirrors upstream fiber.ts:139) + inject.each { |dep, cfg| @ctx = @ctx.intercept(dep, cfg) if cfg.is_a?(Hash) } runtime.fibers << self @ctx.events.emit('internal/plugin', self) inject.each_key { |name| check_impl(name) } @@ -116,7 +121,7 @@ def update(config) # only load/unload execution is async) -- def check_impl(name) - impl = @ctx.root.services[name] + impl = @ctx.root.services[@ctx.isolate_key(name)] if impl && impl.fiber.state == :active @internal_store[name] = impl else diff --git a/spec/cordis/isolate_spec.rb b/spec/cordis/isolate_spec.rb new file mode 100644 index 0000000..cc4a14d --- /dev/null +++ b/spec/cordis/isolate_spec.rb @@ -0,0 +1,127 @@ +# frozen_string_literal: true + +# Mirrors upstream tests/isolate.spec.ts, plus intercept (no upstream standalone spec; +# semantics taken from Service[resolveConfig] and fiber.ts inject-config handling). +RSpec.describe 'isolation and intercept' do + let(:ctx) { Cordis::Context.new } + let(:log) { [] } + + let(:consumer) do + { name: 'consumer', inject: [:foo], apply: lambda { |_c, _config| + log << :load + -> { log << :unload } + } } + end + + describe 'ctx.isolate' do + it 'keeps provides in separate realms per isolated context' do + with_reactor do + ctx.plugin(consumer) + ctx1 = ctx.isolate(:foo) + ctx1.plugin(consumer) + ctx2 = ctx.isolate(:foo) + ctx2.plugin(consumer) + tick + + # provide in the default realm: only the root consumer loads + dispose0 = ctx.provide(:foo, { bar: 100 }) + expect(ctx.get(:foo)).to eq({ bar: 100 }) + expect(ctx1.get(:foo)).to be_nil + expect(ctx2.get(:foo)).to be_nil + tick + expect(log).to eq([:load]) + + # provide in ctx1's realm: only ctx1's consumer loads + ctx1.provide(:foo, { bar: 200 }) + expect(ctx.get(:foo)).to eq({ bar: 100 }) + expect(ctx1.get(:foo)).to eq({ bar: 200 }) + expect(ctx2.get(:foo)).to be_nil + tick + expect(log).to eq(%i[load load]) + + # disposing the default-realm provide only unloads the root consumer + dispose0.call + expect(ctx.get(:foo)).to be_nil + expect(ctx1.get(:foo)).to eq({ bar: 200 }) + tick + expect(log).to eq(%i[load load unload]) + + ctx2.provide(:foo, { bar: 300 }) + expect(ctx2.get(:foo)).to eq({ bar: 300 }) + tick + expect(log).to eq(%i[load load unload load]) + end + end + + it 'shares one realm between isolates created with the same label' do + with_reactor do + label = Object.new + ctx1 = ctx.isolate(:foo, label) + ctx1.plugin(consumer) + ctx2 = ctx.isolate(:foo, label) + ctx2.plugin(consumer) + tick + expect(log).to eq([]) + + ctx.provide(:foo, { bar: 100 }) + tick + expect(log).to eq([]) # default realm is invisible to both + + dispose12 = ctx1.provide(:foo, { bar: 200 }) + expect(ctx1.get(:foo)).to eq({ bar: 200 }) + expect(ctx2.get(:foo)).to eq({ bar: 200 }) + tick + expect(log).to eq(%i[load load]) + + dispose12.call + expect(ctx1.get(:foo)).to be_nil + expect(ctx2.get(:foo)).to be_nil + tick + expect(log).to eq(%i[load load unload unload]) + end + end + + it 'allows the same name to be provided once per realm (duplicate within a realm still raises)' do + with_reactor do + ctx.provide(:foo, 1) + ctx.isolate(:foo).provide(:foo, 2) # different realm: fine + expect { ctx.provide(:foo, 3) }.to raise_error(Cordis::ServiceError, /already been registered/) + end + end + + it 'stops the service walk at an isolation boundary' do + with_reactor do + ctx.provide(:foo, :root_value) + # plugin registered through an isolated ctx must not see the default-realm :foo + fiber = ctx.isolate(:foo).plugin(lambda { |c, _config| + log << begin + c.foo + rescue Cordis::ServiceError => e + e.message + end + }) + fiber.await + expect(log).to eq(['cannot get property "foo" without inject']) + end + end + end + + describe 'ctx.intercept' do + it 'resolves config with inner intercepts overriding outer ones over a base' do + ctx2 = ctx.intercept(:foo, { a: 1, b: 1 }).intercept(:foo, { b: 2 }) + expect(ctx2.resolve_config(:foo, { a: 0, c: 3 })).to eq({ a: 1, b: 2, c: 3 }) + expect(ctx.resolve_config(:foo, { a: 0 })).to eq({ a: 0 }) # original ctx untouched + end + + it 'turns a Hash inject config into an intercept on the plugin ctx' do + with_reactor do + ctx.provide(:foo, :service) + plugin = { name: 'web', inject: { foo: { level: 5 } }, apply: lambda { |c, _config| + log << c.resolve_config(:foo, { level: 1, extra: true }) + } } + ctx.plugin(plugin).await + expect(log).to eq([{ level: 5, extra: true }]) + end + end + end +end