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
5 changes: 5 additions & 0 deletions reflexio/models/api_schema/domain/entities.py
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@
"PlaybookAggregationChangeLog",
"PlaybookAggregationChangeLogResponse",
"OptimizerKind",
"OpenWorldDeploymentLifecycleState",
"OptimizationJobStage",
"OptimizationTerminalOutcome",
"OptimizationArtifactKind",
Expand Down Expand Up @@ -416,6 +417,10 @@ class AgentPlaybook(BaseModel):
"optimizer_legacy_unknown",
]

OpenWorldDeploymentLifecycleState = Literal[
"provisional", "confirmed", "restored", "displaced", "erased"
]

OptimizationJobStage = Literal[
"evidence_frozen",
"discovery_analyzed",
Expand Down
284 changes: 253 additions & 31 deletions reflexio/server/services/playbook/publication.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,14 +9,20 @@
from hashlib import sha256
from typing import Literal, Protocol

from reflexio.models.api_schema.domain.entities import OptimizerKind
from reflexio.models.api_schema.domain.entities import OptimizerKind, UserPlaybook

PublicationOutcome = Literal["applied", "incumbent_changed"]
PublishableOptimizerKind = Literal["gepa", "offline_tuner_replay"]
PublishableOptimizerKind = Literal[
"gepa", "offline_tuner_replay", "offline_tuner_open_world"
]
PublicationSource = Literal["gepa", "offline_optimizer"]

_PUBLISHABLE_OPTIMIZERS = frozenset({"gepa", "offline_tuner_replay"})
_LEGACY_PUBLICATION_OPTIMIZERS = frozenset({"gepa", "offline_tuner_replay"})
_DECISION_PROOF_OPTIMIZERS = frozenset(
{"gepa", "offline_tuner_replay", "offline_tuner_open_world"}
)
_PROJECTION_SCHEMA_VERSION = "offline-tuner-candidate-search-projection-v1"
_USER_PLAYBOOK_FULL_VERSION_SCHEMA = "user-playbook-full-version-v1"
_CANONICAL_DECIMAL = re.compile(r"-?(?:0|[1-9][0-9]*)(?:\.[0-9]*[1-9])?\Z")
PUBLICATION_SUBJECT_EPOCHS_METADATA_KEY = "publication_subject_epochs"
PUBLICATION_PROOF_JSON_METADATA_KEY = "publication_proof_json"
Expand Down Expand Up @@ -106,16 +112,50 @@ def _canonical_payload(name: str, value: str) -> object:
return payload


def _validate_optimizer(value: object) -> None:
if value not in _PUBLISHABLE_OPTIMIZERS:
def _validate_subject_epochs_json(value: str) -> None:
epochs = _canonical_payload("subject_epochs_json", value)
if (
not isinstance(epochs, dict)
or set(epochs) != {"subjects"}
or not isinstance(epochs.get("subjects"), list)
or not epochs["subjects"]
):
raise ValueError("subject epochs must contain a non-empty subjects list")
subject_refs: set[str] = set()
for item in epochs["subjects"]:
if not isinstance(item, dict):
raise ValueError("subject epochs must contain objects")
if set(item) != {"ref", "epoch"}:
raise ValueError("subject epochs must use ref and epoch fields")
subject_ref = item["ref"]
epoch = item["epoch"]
if (
not isinstance(subject_ref, str)
or not subject_ref
or type(epoch) is not int
or epoch < 0
):
raise ValueError("subject epochs contain an invalid identity or epoch")
if subject_ref in subject_refs:
raise ValueError("subject epochs must contain unique subject refs")
subject_refs.add(subject_ref)


def _validate_legacy_publication_optimizer(value: object) -> None:
if value not in _LEGACY_PUBLICATION_OPTIMIZERS:
raise ValueError("optimizer_kind is not publishable")


def _validate_decision_proof_optimizer(value: object) -> None:
if value not in _DECISION_PROOF_OPTIMIZERS:
raise ValueError("optimizer_kind is not publishable")


def publication_source_for_optimizer(
optimizer_kind: OptimizerKind,
) -> PublicationSource:
_validate_optimizer(optimizer_kind)
return "offline_optimizer" if optimizer_kind == "offline_tuner_replay" else "gepa"
_validate_legacy_publication_optimizer(optimizer_kind)
return "offline_optimizer" if optimizer_kind != "gepa" else "gepa"


def incumbent_user_playbook_semantic_digest(
Expand Down Expand Up @@ -156,7 +196,7 @@ class DecisionProofEnvelope:
decision: Literal["apply"]

def __post_init__(self) -> None:
_validate_optimizer(self.optimizer_kind)
_validate_decision_proof_optimizer(self.optimizer_kind)
_require_text("decision proof schema_version", self.schema_version)
_require_digest("decision proof digest", self.digest)
if self.decision != "apply":
Expand Down Expand Up @@ -265,7 +305,7 @@ class PublicationRequest:
request_id: str

def __post_init__(self) -> None:
_validate_optimizer(self.optimizer_kind)
_validate_legacy_publication_optimizer(self.optimizer_kind)
if type(self.job_id) is not int or self.job_id <= 0:
raise ValueError("publication job_id must be positive")
_require_text("publication attempt_key", self.attempt_key)
Expand Down Expand Up @@ -296,32 +336,174 @@ def __post_init__(self) -> None:
self.projection.candidate_content_digest
):
raise ValueError("revised content digest must match search projection")
epochs = _canonical_payload("subject_epochs_json", self.subject_epochs_json)
_validate_subject_epochs_json(self.subject_epochs_json)


@dataclass(frozen=True)
class QualificationAuthorityRef:
epoch: int
authority_digest: str
discovery_component_identity_digest: str
discovery_qualification_suite_digest: str
discovery_qualification_result_digest: str
held_out_component_identity_digest: str
held_out_qualification_suite_digest: str
held_out_qualification_result_digest: str
candidate_generator_identity_digest: str
candidate_generator_authorization_digest: str

def __post_init__(self) -> None:
if type(self.epoch) is not int or self.epoch <= 0:
raise ValueError("qualification authority epoch must be positive")
for field in (
"authority_digest",
"discovery_component_identity_digest",
"discovery_qualification_suite_digest",
"discovery_qualification_result_digest",
"held_out_component_identity_digest",
"held_out_qualification_suite_digest",
"held_out_qualification_result_digest",
"candidate_generator_identity_digest",
"candidate_generator_authorization_digest",
):
_require_digest(f"qualification authority {field}", getattr(self, field))


def _validated_full_version_snapshot(
*,
incumbent_snapshot_json: str,
incumbent_full_version_fingerprint: str,
incumbent_user_playbook_id: int,
) -> UserPlaybook:
_require_digest(
"incumbent full version fingerprint", incumbent_full_version_fingerprint
)
payload = _canonical_payload("incumbent_snapshot_json", incumbent_snapshot_json)
if sha256(incumbent_snapshot_json.encode("utf-8")).hexdigest() != (
incumbent_full_version_fingerprint
):
raise ValueError("incumbent full version fingerprint does not match snapshot")
if (
not isinstance(payload, dict)
or set(payload) != {"schema_version", "user_playbook"}
or payload.get("schema_version") != _USER_PLAYBOOK_FULL_VERSION_SCHEMA
or not isinstance(payload.get("user_playbook"), dict)
):
raise ValueError("incumbent snapshot is not a full user playbook version")
snapshot = UserPlaybook.model_validate(payload["user_playbook"])
expected_playbook = snapshot.model_dump(mode="json", exclude={"embedding"}) | {
"governance_subject_ref": snapshot.governance_subject_ref,
"retired_at": snapshot.retired_at,
}
if canonical_json_bytes(payload["user_playbook"]) != canonical_json_bytes(
expected_playbook
):
raise ValueError("incumbent snapshot is not a full user playbook version")
if snapshot.user_playbook_id != incumbent_user_playbook_id:
raise ValueError("incumbent_user_playbook_id does not match snapshot")
return snapshot


@dataclass(frozen=True)
class ProvisionalPublicationRequest:
optimizer_kind: Literal["offline_tuner_open_world"]
job_id: int
attempt_key: str
publication_claim: PublicationClaim
worker_fence: int
incumbent_user_playbook_id: int
incumbent_full_version_fingerprint: str
incumbent_snapshot_json: str
revised_content: str
projection: PublicationSearchProjection
decision_proof: DecisionProofEnvelope
subject_epochs_json: str
qualification_authority: QualificationAuthorityRef
evidence_bundle_digest: str
candidate_digest: str

def __post_init__(self) -> None:
if self.optimizer_kind != "offline_tuner_open_world":
raise ValueError("optimizer_kind must be offline_tuner_open_world")
if type(self.publication_claim) is not PublicationClaim:
raise ValueError("publication_claim must be PublicationClaim")
if type(self.projection) is not PublicationSearchProjection:
raise ValueError("projection must be PublicationSearchProjection")
if type(self.decision_proof) is not DecisionProofEnvelope:
raise ValueError("decision_proof must be DecisionProofEnvelope")
if type(self.job_id) is not int or self.job_id <= 0:
raise ValueError("provisional publication job_id must be positive")
_require_text("provisional publication attempt_key", self.attempt_key)
if self.publication_claim.job_id != self.job_id:
raise ValueError("publication claim job_id must match request job_id")
if type(self.worker_fence) is not int or self.worker_fence <= 0:
raise ValueError("worker_fence must be positive")
if (
not isinstance(epochs, dict)
or set(epochs) != {"subjects"}
or not isinstance(epochs.get("subjects"), list)
or not epochs["subjects"]
type(self.incumbent_user_playbook_id) is not int
or self.incumbent_user_playbook_id <= 0
):
raise ValueError("incumbent_user_playbook_id must be positive")
snapshot = _validated_full_version_snapshot(
incumbent_snapshot_json=self.incumbent_snapshot_json,
incumbent_full_version_fingerprint=self.incumbent_full_version_fingerprint,
incumbent_user_playbook_id=self.incumbent_user_playbook_id,
)
_require_text("revised_content", self.revised_content)
if self.revised_content == snapshot.content:
raise ValueError("revised_content must differ from incumbent content")
if self.projection.preserved_trigger != snapshot.trigger:
raise ValueError("search projection must preserve incumbent trigger")
if self.decision_proof.optimizer_kind != self.optimizer_kind:
raise ValueError("decision proof optimizer_kind must match request")
if sha256(self.revised_content.encode("utf-8")).hexdigest() != (
self.projection.candidate_content_digest
):
raise ValueError("subject epochs must contain a non-empty subjects list")
subject_refs: set[str] = set()
for item in epochs["subjects"]:
if not isinstance(item, dict):
raise ValueError("subject epochs must contain objects")
if set(item) != {"ref", "epoch"}:
raise ValueError("subject epochs must use ref and epoch fields")
subject_ref = item["ref"]
epoch = item["epoch"]
raise ValueError("revised content digest must match search projection")
_validate_subject_epochs_json(self.subject_epochs_json)
if type(self.qualification_authority) is not QualificationAuthorityRef:
raise ValueError(
"qualification authority must be QualificationAuthorityRef"
)
_require_digest("evidence bundle digest", self.evidence_bundle_digest)
_require_digest("candidate digest", self.candidate_digest)


@dataclass(frozen=True)
class ProvisionalPublicationResult:
job_id: int
outcome: Literal["provisional", "incumbent_changed"]
successor_user_playbook_id: int | None
deployment_lifecycle_id: int | None
observation_deadline: int | None

def __post_init__(self) -> None:
if type(self.job_id) is not int or self.job_id <= 0:
raise ValueError("provisional publication result job_id must be positive")
if self.outcome not in {"provisional", "incumbent_changed"}:
raise ValueError("provisional publication result outcome is invalid")
if self.outcome == "provisional":
if (
not isinstance(subject_ref, str)
or not subject_ref
or type(epoch) is not int
or epoch < 0
type(self.successor_user_playbook_id) is not int
or self.successor_user_playbook_id <= 0
or type(self.deployment_lifecycle_id) is not int
or self.deployment_lifecycle_id <= 0
or type(self.observation_deadline) is not int
or self.observation_deadline <= 0
):
raise ValueError("subject epochs contain an invalid identity or epoch")
if subject_ref in subject_refs:
raise ValueError("subject epochs must contain unique subject refs")
subject_refs.add(subject_ref)
raise ValueError(
"provisional publication requires successor lifecycle and deadline"
)
elif any(
value is not None
for value in (
self.successor_user_playbook_id,
self.deployment_lifecycle_id,
self.observation_deadline,
)
):
raise ValueError(
"incumbent_changed provisional publication cannot have successor state"
)


@dataclass(frozen=True)
Expand Down Expand Up @@ -369,6 +551,41 @@ def load_user_playbook_publication_result(
) -> PublicationResult | None: ...


class UserPlaybookProvisionalPublicationStore(UserPlaybookPublicationStore, Protocol):
"""Durable Phase 4 provisional publication operations."""

def cleanup_user_playbook_provisional_stale_attempt(
self, *, job_id: int, owner: str, worker_fence: int
) -> bool:
"""Clean only matching uncommitted residue; false means ownership changed."""
...

def claim_user_playbook_provisional_publication(
self,
*,
job_id: int,
owner: str,
worker_fence: int,
incumbent_user_playbook_id: int,
incumbent_full_version_fingerprint: str,
incumbent_snapshot_json: str,
) -> PublicationClaim: ...

def stage_user_playbook_provisional_publication(
self, request: ProvisionalPublicationRequest
) -> bool:
"""Stage the request, or return false after cleaning stale-CAS residue."""
...

def commit_user_playbook_provisional_publication(
self, request: ProvisionalPublicationRequest
) -> ProvisionalPublicationResult: ...

def load_user_playbook_provisional_publication_result(
self, job_id: int
) -> ProvisionalPublicationResult | None: ...


class UserPlaybookPublicationService:
"""Coordinates proof verification with durable staging and atomic commit."""

Expand Down Expand Up @@ -418,10 +635,15 @@ def publish_user_playbook_successor(
"PublicationRequest",
"PublicationResult",
"PublicationSearchProjection",
"ProvisionalPublicationRequest",
"ProvisionalPublicationResult",
"PublishableOptimizerKind",
"PUBLICATION_INCUMBENT_CONTENT_DIGEST_METADATA_KEY",
"PUBLICATION_INCUMBENT_SEMANTIC_DIGEST_METADATA_KEY",
"PUBLICATION_INCUMBENT_TRIGGER_METADATA_KEY",
"PUBLICATION_SUBJECT_EPOCHS_METADATA_KEY",
"QualificationAuthorityRef",
"UserPlaybookProvisionalPublicationStore",
"UserPlaybookPublicationService",
"UserPlaybookPublicationStore",
"canonical_json_bytes",
Expand Down
Loading
Loading