mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
fix(#1195): the handlers of one skill still run one at a time
Review finding (gpt-6-astra, triad round 2): the shared-reader barrier let two no-deps callbacks of the SAME skill overlap; the bundled Telegram card renderer reads a message id before awaiting send_message and records it afterwards, so overlapping updates both saw no id and produced duplicate bubbles. The barrier lease now carries the skill as its key: leases of one skill never overlap (the ordering every skill author had under the old exclusive lock), leases of different skills still do (the F3 gain), a deps-bearing writer stays exclusive. Regressions: same-skill sync/async readers serialize, and through the real PluginAPIImpl wrapper one skill's handlers are sequential while another skill's handler runs beside them (fail on the previous tree). Also from this round (gpt-5.6-sol): the record-failure regression test spawns its isolated spawner through platform_layer.subprocess_new_group_kwargs() instead of a raw start_new_session flag. CREATING_SKILLS and the chapter-01 row state the per-skill rule.
This commit is contained in:
parent
6613ea9257
commit
a92dd2e014
6 changed files with 173 additions and 25 deletions
|
|
@ -203,8 +203,10 @@ flowchart LR
|
|||
- **Isolated deps** (pip / npm / uv / node) install into
|
||||
`data/skills/<bucket>/<name>/.ouroboros_env/`. Status is recorded
|
||||
in `data/state/skills/<name>/deps.json`.
|
||||
In-process extension scopes are non-reentrant: no-dependency handlers share
|
||||
a read lease, while a dependency-bearing handler owns an exclusive lease
|
||||
In-process extension scopes are non-reentrant: no-dependency handlers of
|
||||
DIFFERENT skills share a read lease and overlap, the handlers of ONE skill
|
||||
run one at a time (your callbacks stay sequential, as under the old
|
||||
exclusive lock), and a dependency-bearing handler owns an exclusive lease
|
||||
through import, handler waits and cleanup. The async form polls
|
||||
cooperatively, so it never blocks the ASGI loop; nested scopes and
|
||||
unwrapped child work remain unsupported rather than inheriting a lease.
|
||||
|
|
|
|||
|
|
@ -310,7 +310,7 @@ server.py (Starlette+uvicorn) ← HTTP + WebSocket on configurable host:port (de
|
|||
├── extension_process_runner.py ← Extension child processes: scrubbed env, per-skill deps, timeouts, graceful host errors
|
||||
├── extension_route_stream.py ← Portable stdio response frames and ASGI relay for out-of-process extension routes (§3 Out-of-process extension responses)
|
||||
├── extension_ui_validation.py ← The one host-owned declarative-schema-v1 widget validator
|
||||
├── extension_isolated_deps.py ← In-process bridge for isolated-dependency extensions and the `_ExecutionBarrier`: a NON-REENTRANT shared-reader/exclusive-writer lease over the shared `sys.path` seam — no-deps scopes are readers and overlap, a deps-injection scope is a writer that excludes readers across its handler wait and cleanup; sync waits on a Condition, async polls cooperatively so the ASGI loop never blocks
|
||||
├── extension_isolated_deps.py ← In-process bridge for isolated-dependency extensions and the `_ExecutionBarrier`: a NON-REENTRANT shared-reader/exclusive-writer lease over the shared `sys.path` seam — no-deps scopes are readers and overlap ACROSS skills while the scopes of one skill run one at a time (keyed lease: a stateful consumer such as the Telegram card renderer relies on its own callbacks staying sequential), a deps-injection scope is a writer that excludes readers across its handler wait and cleanup; sync waits on a Condition, async polls cooperatively so the ASGI loop never blocks
|
||||
├── extension_health.py ← Durable process-qualified per-skill health at `data/state/skills/<name>/health.json`; server observation is authoritative, worker observation is a handoff-qualified view
|
||||
├── extension_plugin_api.py, extension_registry_state.py, extension_liveness.py, extension_child_catalog.py, extension_import_staging.py, extension_surface_names.py ← The extension runtime's leaves: the `PluginAPI` object handed to `register(api)`; the process-wide registries of live extension surfaces; the liveness authority for one extension; host-side validation of child-catalog surface descriptors; staged import trees and their reclamation; provider-safe naming and syntax rules for extension surfaces
|
||||
├── skill_token.py ← Opaque Host Service token minting/validation (§12)
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@
|
|||
|
||||
Machine extraction of the `docs/ARCHITECTURE.md` "Data layout (`~/Ouroboros/`)" tree — the durable-file orientation carrier (this tree's counterpart of the reference PERSISTENCE_OWNERS derivation checklist) — regenerated by `python scripts/regenerate_inventories.py`. Do not edit. Every entry is probed against reality: repo entries must exist as tracked paths; data-plane entries must appear as a literal in the runtime sources that construct them. A durable file renamed or removed in code while its tree row survives = red (`tests/test_generated_inventories.py`).
|
||||
|
||||
Source: `docs/architecture/01-high-level-architecture.md`, physical LF lines 594-683; UTF-8 SHA-256 `588d1a21240691b69d8c3608bce96fadc310544ce821bf0311d6ca2da5bd852d`.
|
||||
Source: `docs/architecture/01-high-level-architecture.md`, physical LF lines 594-683; UTF-8 SHA-256 `25c79c42ad2b7f8e4a2eca4a96088dc25831779c05360618bdd6dc44e8d41d2d`.
|
||||
|
||||
- entries: **79** (code-ref: 72, repo-dir: 6, repo-path: 1)
|
||||
|
||||
|
|
|
|||
|
|
@ -25,6 +25,13 @@ class _ExecutionBarrier:
|
|||
and blocks new readers once it declares intent. Ownership is deliberately
|
||||
not associated with a thread or task, preserving the old non-reentrant
|
||||
scope contract and avoiding ContextVar inheritance across child tasks.
|
||||
|
||||
A lease may also carry a ``key`` — the skill it executes. Leases of the
|
||||
SAME key never overlap: the old exclusive lock serialized every handler,
|
||||
and a bundled consumer (the Telegram card renderer) relies on its own
|
||||
callbacks running one after another. Independent skills still overlap;
|
||||
only the cross-skill exclusion was the #1195 F3 defect, never the per-skill
|
||||
ordering a skill author may assume.
|
||||
"""
|
||||
|
||||
def __init__(self) -> None:
|
||||
|
|
@ -32,37 +39,51 @@ class _ExecutionBarrier:
|
|||
self.readers = 0
|
||||
self.writer_active = False
|
||||
self.writers_waiting = 0
|
||||
self.active_keys: set[str] = set()
|
||||
|
||||
def try_enter(self, writer: bool) -> bool:
|
||||
def try_enter(self, writer: bool, key: str | None = None) -> bool:
|
||||
with self.condition:
|
||||
if key is not None and key in self.active_keys:
|
||||
return False
|
||||
if writer:
|
||||
if self.writer_active or self.readers:
|
||||
return False
|
||||
self.writer_active = True
|
||||
return True
|
||||
if self.writer_active or self.writers_waiting:
|
||||
return False
|
||||
self.readers += 1
|
||||
else:
|
||||
if self.writer_active or self.writers_waiting:
|
||||
return False
|
||||
self.readers += 1
|
||||
if key is not None:
|
||||
self.active_keys.add(key)
|
||||
return True
|
||||
|
||||
def acquire(self, blocking: bool = True, *, writer: bool = True) -> bool:
|
||||
def acquire(self, blocking: bool = True, *, writer: bool = True, key: str | None = None) -> bool:
|
||||
if not blocking:
|
||||
return self.try_enter(writer)
|
||||
return self.try_enter(writer, key)
|
||||
with self.condition:
|
||||
def key_free() -> bool:
|
||||
return key is None or key not in self.active_keys
|
||||
|
||||
if writer:
|
||||
self.writers_waiting += 1
|
||||
try:
|
||||
self.condition.wait_for(lambda: not self.writer_active and self.readers == 0)
|
||||
self.condition.wait_for(
|
||||
lambda: not self.writer_active and self.readers == 0 and key_free()
|
||||
)
|
||||
self.writer_active = True
|
||||
finally:
|
||||
self.writers_waiting -= 1
|
||||
self.condition.notify_all()
|
||||
else:
|
||||
self.condition.wait_for(lambda: not self.writer_active and not self.writers_waiting)
|
||||
self.condition.wait_for(
|
||||
lambda: not self.writer_active and not self.writers_waiting and key_free()
|
||||
)
|
||||
self.readers += 1
|
||||
if key is not None:
|
||||
self.active_keys.add(key)
|
||||
return True
|
||||
|
||||
def release(self, *, writer: bool = True) -> None:
|
||||
def release(self, *, writer: bool = True, key: str | None = None) -> None:
|
||||
with self.condition:
|
||||
if writer:
|
||||
if not self.writer_active:
|
||||
|
|
@ -72,6 +93,8 @@ class _ExecutionBarrier:
|
|||
if self.readers <= 0:
|
||||
raise RuntimeError("execution reader lease released without enter")
|
||||
self.readers -= 1
|
||||
if key is not None:
|
||||
self.active_keys.discard(key)
|
||||
self.condition.notify_all()
|
||||
|
||||
|
||||
|
|
@ -352,12 +375,20 @@ def _release_site_dirs_best_effort(site_dirs: Sequence[str]) -> None:
|
|||
log.warning("isolated dependency scope cleanup failed after body success: %s", exc)
|
||||
|
||||
|
||||
async def _acquire_execution_barrier_async(*, writer: bool) -> None:
|
||||
def _scope_key(skill_dir: pathlib.Path) -> str:
|
||||
"""One key per skill payload; two spellings of one directory are one skill."""
|
||||
try:
|
||||
return str(pathlib.Path(skill_dir).resolve())
|
||||
except OSError:
|
||||
return str(skill_dir)
|
||||
|
||||
|
||||
async def _acquire_execution_barrier_async(*, writer: bool, key: str | None = None) -> None:
|
||||
if writer:
|
||||
with _execution_lock.condition:
|
||||
_execution_lock.writers_waiting += 1
|
||||
try:
|
||||
while not _execution_lock.try_enter(writer):
|
||||
while not _execution_lock.try_enter(writer, key):
|
||||
await asyncio.sleep(0.01)
|
||||
# No await between grant and the caller's try/finally: cancellation can
|
||||
# only be delivered while polling or inside the protected scope body.
|
||||
|
|
@ -372,13 +403,15 @@ async def _acquire_execution_barrier_async(*, writer: bool) -> None:
|
|||
def isolated_site_dirs_scope(skill_dir: pathlib.Path, *, enabled: bool) -> Iterator[None]:
|
||||
"""Expose reviewed deps in a local reader/writer isolation scope.
|
||||
|
||||
No-deps scopes are readers and overlap. A deps-bearing scope is a writer:
|
||||
it excludes readers while the shared ``sys.path`` is changed and cleaned.
|
||||
Both leases remain non-reentrant and last through handler waits and cleanup.
|
||||
No-deps scopes are readers and overlap ACROSS skills; the scopes of one
|
||||
skill run one at a time. A deps-bearing scope is a writer: it excludes
|
||||
readers while the shared ``sys.path`` is changed and cleaned. Both leases
|
||||
remain non-reentrant and last through handler waits and cleanup.
|
||||
"""
|
||||
|
||||
writer = bool(enabled)
|
||||
_execution_lock.acquire(writer=writer)
|
||||
key = _scope_key(skill_dir)
|
||||
_execution_lock.acquire(writer=writer, key=key)
|
||||
site_dirs: List[str] = []
|
||||
try:
|
||||
site_dirs = inject_isolated_site_dirs(skill_dir) if enabled else []
|
||||
|
|
@ -387,7 +420,7 @@ def isolated_site_dirs_scope(skill_dir: pathlib.Path, *, enabled: bool) -> Itera
|
|||
try:
|
||||
_release_site_dirs_best_effort(site_dirs)
|
||||
finally:
|
||||
_execution_lock.release(writer=writer)
|
||||
_execution_lock.release(writer=writer, key=key)
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
|
|
@ -395,7 +428,8 @@ async def async_isolated_site_dirs_scope(skill_dir: pathlib.Path, *, enabled: bo
|
|||
# Same reader/writer barrier as the sync scope. Async acquisition polls
|
||||
# cooperatively so the ASGI loop never blocks on a threading.Condition.
|
||||
writer = bool(enabled)
|
||||
await _acquire_execution_barrier_async(writer=writer)
|
||||
key = _scope_key(skill_dir)
|
||||
await _acquire_execution_barrier_async(writer=writer, key=key)
|
||||
site_dirs: List[str] = []
|
||||
try:
|
||||
site_dirs = inject_isolated_site_dirs(skill_dir) if enabled else []
|
||||
|
|
@ -404,4 +438,4 @@ async def async_isolated_site_dirs_scope(skill_dir: pathlib.Path, *, enabled: bo
|
|||
try:
|
||||
_release_site_dirs_best_effort(site_dirs)
|
||||
finally:
|
||||
_execution_lock.release(writer=writer)
|
||||
_execution_lock.release(writer=writer, key=key)
|
||||
|
|
|
|||
|
|
@ -40,7 +40,8 @@ def _barrier_is_free_before_and_after():
|
|||
def _free() -> bool:
|
||||
barrier = deps._execution_lock
|
||||
with barrier.condition:
|
||||
return not barrier.writer_active and barrier.readers == 0 and barrier.writers_waiting == 0
|
||||
return (not barrier.writer_active and barrier.readers == 0
|
||||
and barrier.writers_waiting == 0 and not barrier.active_keys)
|
||||
|
||||
|
||||
def _reader(skill_dir, entered: threading.Event, release: threading.Event, sink: list):
|
||||
|
|
@ -661,3 +662,111 @@ def test_a_deps_bearing_load_excludes_a_concurrent_no_deps_load(tmp_path):
|
|||
extension_loader.unload_extension("reader_neighbour")
|
||||
finally:
|
||||
del builtins._ouro_1195_gate
|
||||
|
||||
|
||||
# One skill, one handler at a time: the per-skill sequencing a consumer such as
|
||||
# the bundled Telegram card renderer relies on (read the message id, await the
|
||||
# send, record it). Cross-skill overlap is the F3 gain; same-skill ordering is
|
||||
# the contract the old exclusive lock gave every skill author for free.
|
||||
|
||||
|
||||
def test_same_skill_sync_readers_run_one_at_a_time(tmp_path):
|
||||
skill_dir = tmp_path / "telegram"
|
||||
first_in, first_out, second_in = threading.Event(), threading.Event(), threading.Event()
|
||||
sink: list = []
|
||||
first = threading.Thread(target=_reader(skill_dir, first_in, first_out, sink))
|
||||
second_release = threading.Event()
|
||||
second = threading.Thread(target=_reader(skill_dir, second_in, second_release, sink))
|
||||
first.start()
|
||||
assert first_in.wait(GATE)
|
||||
second.start()
|
||||
assert not second_in.wait(0.3), "a second handler of the SAME skill entered while the first was live"
|
||||
assert not deps._execution_lock.try_enter(False, deps._scope_key(skill_dir))
|
||||
assert deps._execution_lock.try_enter(False, deps._scope_key(tmp_path / "other"))
|
||||
deps._execution_lock.release(writer=False, key=deps._scope_key(tmp_path / "other"))
|
||||
first_out.set()
|
||||
assert second_in.wait(GATE), "the second handler never entered after the first left"
|
||||
second_release.set()
|
||||
first.join(GATE)
|
||||
second.join(GATE)
|
||||
assert sink == ["reader", "reader"], sink
|
||||
|
||||
|
||||
def test_same_skill_async_handlers_run_one_at_a_time(tmp_path):
|
||||
async def main():
|
||||
skill_dir = tmp_path / "telegram"
|
||||
first_in = asyncio.Event()
|
||||
first_out = asyncio.Event()
|
||||
order: list = []
|
||||
|
||||
async def handler(label, gate_in, gate_out):
|
||||
async with deps.async_isolated_site_dirs_scope(skill_dir, enabled=False):
|
||||
order.append(f"{label}:in")
|
||||
if gate_in is not None:
|
||||
gate_in.set()
|
||||
if gate_out is not None:
|
||||
await gate_out.wait()
|
||||
order.append(f"{label}:out")
|
||||
|
||||
one = asyncio.create_task(handler("one", first_in, first_out))
|
||||
await first_in.wait()
|
||||
two = asyncio.create_task(handler("two", None, None))
|
||||
await asyncio.sleep(0.2)
|
||||
assert order == ["one:in"], order # two is waiting, not inside
|
||||
first_out.set()
|
||||
await asyncio.wait_for(asyncio.gather(one, two), GATE)
|
||||
assert order == ["one:in", "one:out", "two:in", "two:out"], order
|
||||
|
||||
asyncio.run(main())
|
||||
|
||||
|
||||
def test_registered_handlers_of_one_skill_serialize_and_of_two_skills_overlap(tmp_path):
|
||||
"""Through the real PluginAPIImpl wrapper: same skill sequential, other skill concurrent."""
|
||||
from ouroboros.extension_plugin_api import PluginAPIImpl, _PluginAPIConfig
|
||||
|
||||
def api_for(label):
|
||||
state_dir = tmp_path / "state" / label
|
||||
state_dir.mkdir(parents=True, exist_ok=True)
|
||||
return PluginAPIImpl(_PluginAPIConfig(
|
||||
skill_name=label, permissions=[], env_allowlist=[], state_dir=state_dir,
|
||||
settings_reader=lambda: {}, skill_dir=tmp_path / label,
|
||||
))
|
||||
|
||||
api_a = api_for("skill_a")
|
||||
api_b = api_for("skill_b")
|
||||
|
||||
a_inside = threading.Event()
|
||||
a_release = threading.Event()
|
||||
seen: list = []
|
||||
|
||||
def slow_a(*_args):
|
||||
seen.append("a:in")
|
||||
a_inside.set()
|
||||
a_release.wait(GATE)
|
||||
seen.append("a:out")
|
||||
|
||||
def quick(label):
|
||||
def run(*_args):
|
||||
seen.append(f"{label}:in")
|
||||
seen.append(f"{label}:out")
|
||||
return run
|
||||
|
||||
wrapped_slow_a = api_a._wrap_runtime_handler(slow_a)
|
||||
wrapped_quick_a = api_a._wrap_runtime_handler(quick("a2"))
|
||||
wrapped_quick_b = api_b._wrap_runtime_handler(quick("b"))
|
||||
|
||||
slow = threading.Thread(target=wrapped_slow_a)
|
||||
slow.start()
|
||||
assert a_inside.wait(GATE)
|
||||
other = threading.Thread(target=wrapped_quick_b)
|
||||
other.start()
|
||||
other.join(GATE)
|
||||
assert not other.is_alive(), "another skill's handler was blocked by skill_a's live handler"
|
||||
assert seen[:3] == ["a:in", "b:in", "b:out"], seen
|
||||
same = threading.Thread(target=wrapped_quick_a)
|
||||
same.start()
|
||||
assert not same.join(0.3) and same.is_alive(), "skill_a's second handler overlapped its first"
|
||||
a_release.set()
|
||||
same.join(GATE)
|
||||
slow.join(GATE)
|
||||
assert seen == ["a:in", "b:in", "b:out", "a:out", "a2:in", "a2:out"], seen
|
||||
|
|
|
|||
|
|
@ -12,6 +12,8 @@ import sys
|
|||
|
||||
import pytest
|
||||
|
||||
from ouroboros.platform_layer import subprocess_new_group_kwargs
|
||||
|
||||
REPO_ROOT = pathlib.Path(__file__).resolve().parents[1]
|
||||
_POSIX_ONLY = pytest.mark.skipif(os.name == "nt", reason="POSIX process groups only")
|
||||
|
||||
|
|
@ -49,7 +51,8 @@ def test_record_failure_kills_only_the_child_never_the_spawners_group(tmp_path):
|
|||
spawner = subprocess.run(
|
||||
[sys.executable, "-c", _RECORD_FAILURE_SPAWNER, str(tmp_path)],
|
||||
cwd=str(REPO_ROOT), capture_output=True, text=True, timeout=60,
|
||||
start_new_session=True, env={**os.environ, "PYTHONPATH": str(REPO_ROOT)},
|
||||
env={**os.environ, "PYTHONPATH": str(REPO_ROOT)},
|
||||
**subprocess_new_group_kwargs(),
|
||||
)
|
||||
assert spawner.returncode == 0, spawner.stderr
|
||||
assert spawner.stdout.split() == ["survived", "OSError"], (spawner.stdout, spawner.stderr)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue