Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -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
Expand All @@ -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 — 骨架穩了再說

Expand Down Expand Up @@ -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
Expand All @@ -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 — 骨格が安定してから

Expand Down
1 change: 1 addition & 0 deletions agents.md
Original file line number Diff line number Diff line change
Expand Up @@ -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。

## 範圍紀律
Expand Down
18 changes: 18 additions & 0 deletions examples/demo.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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}"
Expand Down
79 changes: 61 additions & 18 deletions lib/cordis/context.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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, &)
Expand Down Expand Up @@ -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
Expand All @@ -79,36 +114,38 @@ 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
end

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)

impl.value = 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

Expand All @@ -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
Expand Down
11 changes: 8 additions & 3 deletions lib/cordis/fiber.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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) }
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading