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
4 changes: 4 additions & 0 deletions loopx/control_plane/coordination/shadow_management.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,10 @@ def shadow_management_directory(runtime_root: Path, goal_id: str) -> Path:
return runtime_root / "authority-transition" / "file-v0" / f"shadow-management-{digest}"


def runtime_artifact_lock_target(runtime_root: Path, goal_id: str) -> Path:
return shadow_management_directory(runtime_root, goal_id) / "runtime-artifacts"


def shadow_maintenance_lock_target(runtime_root: Path, goal_id: str) -> Path:
return shadow_management_directory(runtime_root, goal_id) / "maintenance"

Expand Down
71 changes: 54 additions & 17 deletions loopx/control_plane/goals/first_party_host_admission.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
cross_runtime_lock_witness,
exclusive_cross_runtime_file_lock,
)
from ..coordination.shadow_management import runtime_artifact_lock_target
from ..effect_runtime import effect_runtime_result
from ..projects.registry_codec import (
SOURCE_SESSION_PROFILE_ID,
Expand Down Expand Up @@ -68,20 +69,44 @@ def _effect(
journal_path=self.journal_path,
)

@contextmanager
def _runtime_artifact_lifetime(self, *, operation: str) -> Iterator[None]:
path = self.journal_path.expanduser().resolve()
turns_dir = path.parent
goal_dir = turns_dir.parent
goals_dir = goal_dir.parent
if (
turns_dir.name != "turns"
or goal_dir.name != self.goal_admission.goal_id
or goals_dir.name != "goals"
):
raise FirstPartyHostRuntimeRejected("journal_path_mismatch")
with exclusive_cross_runtime_file_lock(
runtime_artifact_lock_target(
goals_dir.parent,
self.goal_admission.goal_id,
),
operation=operation,
):
yield

def prepare(
self,
step_kind: TurnProviderStepKind,
effect_ref: str,
persist_journal: JournalPersist,
) -> None:
try:
prepare_source_turn_effect(
registry_path=self.goal_admission.registry_path,
goal_id=self.goal_admission.goal_id,
effect=self._effect(step_kind, effect_ref),
source_admission=self.goal_admission.source_journal_admission_locked,
persist_journal=persist_journal,
)
with self._runtime_artifact_lifetime(
operation="source_turn_effect_runtime_artifact_admit",
):
prepare_source_turn_effect(
registry_path=self.goal_admission.registry_path,
goal_id=self.goal_admission.goal_id,
effect=self._effect(step_kind, effect_ref),
source_admission=self.goal_admission.source_journal_admission_locked,
persist_journal=persist_journal,
)
except SourceTurnEffectRejected as exc:
raise FirstPartyHostRuntimeRejected(exc.code) from exc

Expand All @@ -100,13 +125,16 @@ def release(
persist_journal: JournalPersist,
) -> None:
try:
release_source_turn_effect(
registry_path=self.goal_admission.registry_path,
goal_id=self.goal_admission.goal_id,
effect=self._effect(step_kind, effect_ref),
source_admission=self.goal_admission.source_journal_admission_locked,
persist_journal=persist_journal,
)
with self._runtime_artifact_lifetime(
operation="source_turn_effect_runtime_artifact_release",
):
release_source_turn_effect(
registry_path=self.goal_admission.registry_path,
goal_id=self.goal_admission.goal_id,
effect=self._effect(step_kind, effect_ref),
source_admission=self.goal_admission.source_journal_admission_locked,
persist_journal=persist_journal,
)
except SourceTurnEffectRejected as exc:
raise FirstPartyHostRuntimeRejected(exc.code) from exc

Expand Down Expand Up @@ -310,6 +338,8 @@ def current_lifetime(self, *, operation: str) -> Iterator[None]:
@contextmanager
def source_journal_admission(
self,
*,
runtime_root: Path,
) -> Iterator[dict[str, Any] | None]:
"""Hand one journal mutation to the TS owner under the source guard."""

Expand All @@ -318,10 +348,17 @@ def source_journal_admission(
return
target = guard_path(self.registry_path, self.goal_id)
with exclusive_cross_runtime_file_lock(
target,
operation="first_party_host_journal_commit",
runtime_artifact_lock_target(
runtime_root.expanduser().resolve(),
self.goal_id,
),
operation="first_party_host_journal_runtime_artifact_guard",
):
yield self.source_journal_admission_locked()
with exclusive_cross_runtime_file_lock(
target,
operation="first_party_host_journal_commit",
):
yield self.source_journal_admission_locked()

def source_journal_admission_locked(self) -> dict[str, Any]:
"""Build a TS handoff while the caller holds this Goal's source guard."""
Expand Down
84 changes: 55 additions & 29 deletions loopx/control_plane/quota/spend_commit.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,11 @@
from pathlib import Path
from typing import Any

from ...file_lock import exclusive_file_lock
from ...file_lock import exclusive_cross_runtime_file_lock
from ..coordination.shadow_management import (
runtime_artifact_lock_target,
shadow_maintenance_lock_target,
)
from ..effect_runtime import (
EffectRuntimeConflict,
EffectRuntimeRejected,
Expand Down Expand Up @@ -226,38 +230,60 @@ def record_quota_slot_spend_from_preview(
):
raise ValueError("quota spend preview index basis must be a string or null")
if execute:
index_path = runtime_root / "goals" / safe_goal_id / "runs" / "index.jsonl"
if registry_path is None and goal_ref is None:
with exclusive_file_lock(
index_path,
operation="quota_spend_commit",
):
result = _quota_spend_commit_result(
preview,
source=source,
generated_at=None,
execute=True,
runtime_root=runtime_root,
expected_index_digest=expected_index_digest,
)
else:
with quota_accounting_admission(
resolved_runtime_root = runtime_root.resolve()
artifact_target = runtime_artifact_lock_target(
resolved_runtime_root,
safe_goal_id,
)
maintenance_target = shadow_maintenance_lock_target(
resolved_runtime_root,
safe_goal_id,
)

def commit(source_admission: Mapping[str, Any] | None) -> Mapping[str, Any]:
return _quota_spend_commit_result(
preview,
source=source,
generated_at=None,
execute=True,
runtime_root=runtime_root,
registry_path=registry_path,
goal_id=safe_goal_id,
expected_index_digest=expected_index_digest,
goal_ref=goal_ref,
operation="quota_spend_commit",
) as source_admission:
result = _quota_spend_commit_result(
preview,
source=source,
generated_at=None,
execute=True,
source_admission=source_admission,
)

with exclusive_cross_runtime_file_lock(
artifact_target,
operation="quota_spend_runtime_artifact_guard",
):
# Legacy commits have no source-lifetime guard, so maintenance must
# precede their index lock. Exact commits retain index -> source -> M.
if goal_ref is None:
with exclusive_cross_runtime_file_lock(
maintenance_target,
operation="quota_spend_runtime_artifact_commit",
):
with quota_accounting_admission(
runtime_root=runtime_root,
registry_path=registry_path,
goal_id=safe_goal_id,
goal_ref=None,
operation="quota_spend_commit",
) as source_admission:
result = commit(source_admission)
else:
with quota_accounting_admission(
runtime_root=runtime_root,
expected_index_digest=expected_index_digest,
registry_path=registry_path,
goal_id=safe_goal_id,
goal_ref=goal_ref,
source_admission=source_admission,
)
operation="quota_spend_commit",
) as source_admission:
with exclusive_cross_runtime_file_lock(
maintenance_target,
operation="quota_spend_runtime_artifact_commit",
):
result = commit(source_admission)
else:
result = _quota_spend_commit_result(
preview,
Expand Down
4 changes: 3 additions & 1 deletion loopx/control_plane/turn_driver/executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -1400,7 +1400,9 @@ def persist_journal(snapshot: Mapping[str, Any]) -> None:
_write_journal(journal_path, snapshot)
return
try:
with goal_admission.source_journal_admission() as source_admission:
with goal_admission.source_journal_admission(
runtime_root=runtime_root,
) as source_admission:
if source_admission is None:
raise RuntimeError("source journal admission was not produced")
_write_journal(
Expand Down
51 changes: 44 additions & 7 deletions loopx/control_plane/turn_driver/journal_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,11 @@
from pathlib import Path
from typing import Any

from ...file_lock import exclusive_file_lock
from ...file_lock import exclusive_cross_runtime_file_lock, exclusive_file_lock
from ..coordination.shadow_management import (
runtime_artifact_lock_target,
shadow_maintenance_lock_target,
)
from ..effect_runtime import EffectRuntimeConflict, effect_runtime_result
from .turn_journal_runtime import (
write_turn_journal,
Expand Down Expand Up @@ -57,18 +61,51 @@ def journal_committed_effect_id(journal: Mapping[str, Any]) -> str | None:
return effect_id or None


def _turn_journal_goal_route(path: Path) -> tuple[Path, str] | None:
resolved = path.expanduser().resolve()
turns_dir = resolved.parent
goal_dir = turns_dir.parent
goals_dir = goal_dir.parent
if turns_dir.name != "turns" or goals_dir.name != "goals":
return None
return goals_dir.parent, goal_dir.name


def write_turn_journal_checkpoint(
path: Path,
journal: Mapping[str, Any],
*,
source_admission: Mapping[str, Any] | None = None,
) -> None:
write_turn_journal(
str(path),
journal,
expected_effect_id=journal_committed_effect_id(journal),
source_admission=source_admission,
)
def write() -> None:
write_turn_journal(
str(path),
journal,
expected_effect_id=journal_committed_effect_id(journal),
source_admission=source_admission,
)

route = _turn_journal_goal_route(path)
if route is None:
write()
return
runtime_root, goal_id = route

def write_under_maintenance() -> None:
with exclusive_cross_runtime_file_lock(
shadow_maintenance_lock_target(runtime_root, goal_id),
operation="turn_journal_runtime_artifact_commit",
):
write()

if source_admission is not None:
write_under_maintenance()
return
with exclusive_cross_runtime_file_lock(
runtime_artifact_lock_target(runtime_root, goal_id),
operation="turn_journal_runtime_artifact_guard",
):
write_under_maintenance()


def load_loopx_turn_plan_from_journal(
Expand Down
Loading
Loading