Skip to content

Generated API reference

These signatures and docstrings are extracted from the source in this checkout. Use the public API map for supported import paths and the integration guides for credentials and account tracks.

Capability contracts

zeo_core.contracts.CapabilityResult

Bases: BaseModel, Generic[T]

Standard return envelope for ALL capabilities.

Orchestrators (n8n, Temporal) parse this JSON to decide the next step in the workflow. This enables: - Machine branching (success/skip/error paths) - Audit trails (logs, timing, metadata) - Debugging (structured errors with context)

Invariants
  • If status == error, then error must be present AND machine_message must be present
  • If status == error, then machine_message must start with a recognized prefix (ZEO_ preferred; QC_/ZC_ accepted, see below)
  • If status == success, then error must be None AND machine_message should be None
  • If status == skipped, then error must be None AND machine_message must be present
  • If status == skipped, then machine_message must start with a recognized prefix (ZEO_ preferred; QC_/ZC_ accepted, see below)
  • If machine_message is present, it must start with a recognized prefix (ZEO_ preferred; QC_/ZC_ accepted, see below)

Machine-message / error-code prefix: New code should use ZEO_<AREA>_<DETAIL> (matches the package name, zeo_core). QC_<AREA>_<DETAIL> is still accepted -- it is the convention this package's capabilities used before the pre-extraction rename from quack_core to zeo_core, and it remains valid on purpose (not a leftover bug) because orchestrators already branch on QC_* codes emitted by existing tools; widening the validator to also accept ZEO_/ZC_ was chosen over a hard rename so no existing call site or downstream consumer breaks. ZC_ is accepted as a short-form alias of ZEO_.

Usage Pattern

Tools should use the helper methods (.ok(), .skip(), .fail()) rather than constructing CapabilityResult directly:

Success

result = CapabilityResult.ok( ... data={"transcription": "Hello world"}, ... msg="Transcription completed" ... )

Skip

result = CapabilityResult.skip( ... reason="Video too short for processing", ... code="ZEO_VAL_TOO_SHORT" ... )

Error

try: ... risky_operation() ... except Exception as e: ... result = CapabilityResult.fail_from_exc( ... msg="Failed to process video", ... code="ZEO_IO_ERROR", ... exc=e ... )

Source code in src/zeo_core/contracts/envelopes/result.py
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
class CapabilityResult(BaseModel, Generic[T]):
    """
    Standard return envelope for ALL capabilities.

    Orchestrators (n8n, Temporal) parse this JSON to decide the next step
    in the workflow. This enables:
    - Machine branching (success/skip/error paths)
    - Audit trails (logs, timing, metadata)
    - Debugging (structured errors with context)

    Invariants:
        - If status == error, then error must be present AND machine_message
          must be present
        - If status == error, then machine_message must start with a
          recognized prefix (ZEO_ preferred; QC_/ZC_ accepted, see below)
        - If status == success, then error must be None AND machine_message
          should be None
        - If status == skipped, then error must be None AND machine_message
          must be present
        - If status == skipped, then machine_message must start with a
          recognized prefix (ZEO_ preferred; QC_/ZC_ accepted, see below)
        - If machine_message is present, it must start with a recognized
          prefix (ZEO_ preferred; QC_/ZC_ accepted, see below)

    Machine-message / error-code prefix:
        New code should use ``ZEO_<AREA>_<DETAIL>`` (matches the package
        name, zeo_core). ``QC_<AREA>_<DETAIL>`` is still accepted -- it is
        the convention this package's capabilities used before the
        pre-extraction rename from ``quack_core`` to ``zeo_core``, and it
        remains valid on purpose (not a leftover bug) because orchestrators
        already branch on `QC_*` codes emitted by existing tools; widening
        the validator to also accept `ZEO_`/`ZC_` was chosen over a hard
        rename so no existing call site or downstream consumer breaks.
        ``ZC_`` is accepted as a short-form alias of ``ZEO_``.

    Usage Pattern:
        Tools should use the helper methods (.ok(), .skip(), .fail()) rather
        than constructing CapabilityResult directly:

        >>> # Success
        >>> result = CapabilityResult.ok(
        ...     data={"transcription": "Hello world"},
        ...     msg="Transcription completed"
        ... )

        >>> # Skip
        >>> result = CapabilityResult.skip(
        ...     reason="Video too short for processing",
        ...     code="ZEO_VAL_TOO_SHORT"
        ... )

        >>> # Error
        >>> try:
        ...     risky_operation()
        ... except Exception as e:
        ...     result = CapabilityResult.fail_from_exc(
        ...         msg="Failed to process video",
        ...         code="ZEO_IO_ERROR",
        ...         exc=e
        ...     )
    """

    model_config = ConfigDict(
        extra="forbid",  # Strict schema - no unexpected fields
    )

    # Core status
    status: CapabilityStatus = Field(
        ..., description="Execution status for machine branching"
    )

    outcome: CapabilityOutcome | None = Field(
        None,
        description=(
            "Fine-grained outcome. Defaults from status for legacy constructors: "
            "success→success, skipped→policy_skipped, error→integration_failure. "
            "The invoke helper always sets this explicitly."
        ),
    )

    # Payload (the actual value produced by the capability)
    data: T | None = Field(
        None, description="The actual result data (type varies by capability)"
    )

    # Telemetry
    run_id: str = Field(
        default_factory=generate_run_id,
        description=(
            "Unique identifier for this execution (should match RunManifest.run_id)"
        ),
    )

    timestamp: datetime = Field(
        default_factory=utcnow, description="UTC timestamp when result was created"
    )

    duration_sec: float | None = Field(
        None, ge=0.0, description="Execution duration in seconds (None if not measured)"
    )

    # Messages
    human_message: str = Field(..., description="Readable summary for logs/CLI/UI")

    machine_message: str | None = Field(
        None,
        description=(
            "Machine-readable code for orchestrator branching "
            "(must start with ZEO_, ZC_, or the legacy QC_)"
        ),
    )

    # Diagnostics
    error: CapabilityError | None = Field(
        None, description="Structured error info if status == error"
    )

    logs: list[CapabilityLogEvent] = Field(
        default_factory=list, description="Structured log events from execution"
    )

    metadata: dict[str, Any] = Field(
        default_factory=dict,
        description="Additional context (tool, version, config, etc.)",
    )

    #: Recognized machine-message / error-code prefixes. ZEO_ is the
    #: preferred, current convention (matches the zeo_core package name);
    #: QC_ is accepted for backward compatibility with call sites and
    #: orchestrators that predate the quack_core -> zeo_core rename; ZC_ is
    #: accepted as ZEO_'s short-form alias. See CapabilityResult's class
    #: docstring for the full rationale.
    MACHINE_MESSAGE_PREFIXES: ClassVar[tuple[str, ...]] = ("ZEO_", "ZC_", "QC_")

    @field_validator("machine_message")
    @classmethod
    def validate_machine_message_format(cls, v: str | None) -> str | None:
        """Ensure machine_message follows a recognized *_ convention when present."""
        if v is not None and not v.startswith(cls.MACHINE_MESSAGE_PREFIXES):
            raise ValueError(
                f"machine_message must start with one of "
                f"{cls.MACHINE_MESSAGE_PREFIXES}, got: {v}. "
                "Use format: ZEO_<AREA>_<DETAIL> (e.g., ZEO_VAL_TOO_SHORT)"
            )
        return v

    @model_validator(mode="after")
    def validate_status_invariants(self) -> "CapabilityResult[T]":
        """
        Enforce invariants between status and other fields.

        This ensures orchestrators can rely on the structure:
        - Errors always have error objects and machine codes
        - Successes never have error objects or machine codes
        - Skips never have error objects but always have machine codes
        """
        if self.status == CapabilityStatus.error:
            if self.error is None:
                raise ValueError("status=error requires error field to be present")
            if self.machine_message is None:
                raise ValueError("status=error requires machine_message for branching")

        if self.status == CapabilityStatus.success:
            if self.error is not None:
                raise ValueError("status=success must not have error field")
            if self.machine_message is not None:
                raise ValueError(
                    "status=success should not have machine_message "
                    "(success is the default path, no special routing needed)"
                )

        if self.status == CapabilityStatus.skipped:
            if self.error is not None:
                raise ValueError(
                    "status=skipped must not have error field "
                    "(skips are policy decisions, not errors)"
                )
            if self.machine_message is None:
                raise ValueError(
                    "status=skipped requires machine_message for branching"
                )

        self._coerce_outcome()
        return self

    def _coerce_outcome(self) -> None:
        if self.outcome is None:
            if self.status == CapabilityStatus.success:
                self.outcome = CapabilityOutcome.success
            elif self.status == CapabilityStatus.skipped:
                self.outcome = CapabilityOutcome.policy_skipped
            else:
                self.outcome = CapabilityOutcome.integration_failure
        elif OUTCOME_TO_STATUS[self.outcome] != self.status:
            raise ValueError(
                f"outcome {self.outcome} is incompatible with status {self.status}"
            )

    # Convenience constructors

    @classmethod
    def ok(
        cls,
        data: T,
        msg: str = "Success",
        metadata: dict[str, Any] | None = None,
        logs: list[CapabilityLogEvent] | None = None,
        duration_sec: float | None = None,
        run_id: str | None = None,
    ) -> "CapabilityResult[T]":
        """
        Create a successful result.

        Args:
            data: The result payload
            msg: Human-readable success message
            metadata: Optional metadata dict
            logs: Optional log events from execution
            duration_sec: Execution time in seconds (None if not measured)
            run_id: Optional run_id to reuse (should match manifest run_id)

        Returns:
            CapabilityResult with status=success

        Example:
            >>> result = CapabilityResult.ok(
            ...     data={"clips": [...]},
            ...     msg="Generated 5 clips",
            ...     metadata={"tool": "slice_video", "preset": "fast"}
            ... )
        """
        kwargs = {
            "status": CapabilityStatus.success,
            "outcome": CapabilityOutcome.success,
            "data": data,
            "human_message": msg,
            "metadata": metadata or {},
            "logs": logs or [],
            "duration_sec": duration_sec,
        }
        if run_id is not None:
            kwargs["run_id"] = run_id
        return cls(**kwargs)

    @classmethod
    def skip(
        cls,
        reason: str,
        code: str,
        metadata: dict[str, Any] | None = None,
        run_id: str | None = None,
    ) -> "CapabilityResult[T]":
        """
        Create a skip result (valid policy decision).

        Skips are NOT errors - they represent intentional decisions
        to skip processing (e.g., video too short, file already exists).

        Args:
            reason: Human-readable explanation for the skip
            code: Machine-readable skip code (must start with ZEO_, ZC_, or
                the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)
            metadata: Optional metadata dict
            run_id: Optional run_id to reuse (should match manifest run_id)

        Returns:
            CapabilityResult with status=skipped

        Example:
            >>> result = CapabilityResult.skip(
            ...     reason="Video duration under 10 seconds",
            ...     code="ZEO_VAL_TOO_SHORT"
            ... )
        """
        kwargs = {
            "status": CapabilityStatus.skipped,
            "outcome": CapabilityOutcome.policy_skipped,
            "human_message": reason,
            "machine_message": code,
            "metadata": metadata or {},
        }
        if run_id is not None:
            kwargs["run_id"] = run_id
        return cls(**kwargs)

    @classmethod
    def unavailable(
        cls,
        reason: str,
        code: str = "ZEO_CAP_UNAVAILABLE",
        metadata: dict[str, Any] | None = None,
        run_id: str | None = None,
    ) -> "CapabilityResult[T]":
        """Create a skipped result because a declared dependency is missing."""
        kwargs: dict[str, Any] = {
            "status": CapabilityStatus.skipped,
            "outcome": CapabilityOutcome.unavailable,
            "human_message": reason,
            "machine_message": code,
            "metadata": metadata or {},
        }
        if run_id is not None:
            kwargs["run_id"] = run_id
        return cls(**kwargs)

    @classmethod
    def fail(
        cls,
        msg: str,
        code: str,
        exception: Exception | None = None,
        metadata: dict[str, Any] | None = None,
        logs: list[CapabilityLogEvent] | None = None,
        run_id: str | None = None,
        outcome: CapabilityOutcome = CapabilityOutcome.integration_failure,
    ) -> "CapabilityResult[T]":
        """
        Create an error result.

        Args:
            msg: Human-readable error message
            code: Machine-readable error code (must start with ZEO_, ZC_, or
                the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)
            exception: Optional exception that caused the error
            metadata: Optional metadata dict
            logs: Optional log events from execution
            run_id: Optional run_id to reuse (should match manifest run_id)

        Returns:
            CapabilityResult with status=error

        Example:
            >>> result = CapabilityResult.fail(
            ...     msg="Failed to read video file",
            ...     code="ZEO_IO_NOT_FOUND",
            ...     exception=FileNotFoundError("/data/video.mp4")
            ... )
        """
        err_details: dict[str, Any] = {}
        if exception:
            err_details = {
                "type": type(exception).__name__,
                "str": str(exception),
            }

        kwargs = {
            "status": CapabilityStatus.error,
            "outcome": outcome,
            "human_message": msg,
            "machine_message": code,
            "error": CapabilityError(code=code, message=msg, details=err_details),
            "metadata": metadata or {},
            "logs": logs or [],
        }
        if run_id is not None:
            kwargs["run_id"] = run_id
        return cls(**kwargs)

    @classmethod
    def fail_from_exc(
        cls,
        msg: str,
        code: str,
        exc: Exception,
        metadata: dict[str, Any] | None = None,
        run_id: str | None = None,
    ) -> "CapabilityResult[T]":
        """
        Convenience wrapper for fail() that always includes exception.

        Args:
            msg: Human-readable error message
            code: Machine-readable error code (must start with ZEO_, ZC_, or
                the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)
            exc: Exception that caused the error
            metadata: Optional metadata dict
            run_id: Optional run_id to reuse (should match manifest run_id)

        Returns:
            CapabilityResult with status=error

        Example:
            >>> try:
            ...     process_video()
            ... except IOError as e:
            ...     result = CapabilityResult.fail_from_exc(
            ...         msg="Video processing failed",
            ...         code="ZEO_IO_ERROR",
            ...         exc=e
            ...     )
        """
        return cls.fail(
            msg=msg,
            code=code,
            exception=exc,
            metadata=metadata,
            run_id=run_id,
            outcome=CapabilityOutcome.unexpected_exception,
        )

fail classmethod

fail(msg, code, exception=None, metadata=None, logs=None, run_id=None, outcome=CapabilityOutcome.integration_failure)

Create an error result.

Parameters:

Name Type Description Default
msg str

Human-readable error message

required
code str

Machine-readable error code (must start with ZEO_, ZC_, or the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)

required
exception Exception | None

Optional exception that caused the error

None
metadata dict[str, Any] | None

Optional metadata dict

None
logs list[CapabilityLogEvent] | None

Optional log events from execution

None
run_id str | None

Optional run_id to reuse (should match manifest run_id)

None

Returns:

Type Description
CapabilityResult[T]

CapabilityResult with status=error

Example

result = CapabilityResult.fail( ... msg="Failed to read video file", ... code="ZEO_IO_NOT_FOUND", ... exception=FileNotFoundError("/data/video.mp4") ... )

Source code in src/zeo_core/contracts/envelopes/result.py
@classmethod
def fail(
    cls,
    msg: str,
    code: str,
    exception: Exception | None = None,
    metadata: dict[str, Any] | None = None,
    logs: list[CapabilityLogEvent] | None = None,
    run_id: str | None = None,
    outcome: CapabilityOutcome = CapabilityOutcome.integration_failure,
) -> "CapabilityResult[T]":
    """
    Create an error result.

    Args:
        msg: Human-readable error message
        code: Machine-readable error code (must start with ZEO_, ZC_, or
            the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)
        exception: Optional exception that caused the error
        metadata: Optional metadata dict
        logs: Optional log events from execution
        run_id: Optional run_id to reuse (should match manifest run_id)

    Returns:
        CapabilityResult with status=error

    Example:
        >>> result = CapabilityResult.fail(
        ...     msg="Failed to read video file",
        ...     code="ZEO_IO_NOT_FOUND",
        ...     exception=FileNotFoundError("/data/video.mp4")
        ... )
    """
    err_details: dict[str, Any] = {}
    if exception:
        err_details = {
            "type": type(exception).__name__,
            "str": str(exception),
        }

    kwargs = {
        "status": CapabilityStatus.error,
        "outcome": outcome,
        "human_message": msg,
        "machine_message": code,
        "error": CapabilityError(code=code, message=msg, details=err_details),
        "metadata": metadata or {},
        "logs": logs or [],
    }
    if run_id is not None:
        kwargs["run_id"] = run_id
    return cls(**kwargs)

fail_from_exc classmethod

fail_from_exc(msg, code, exc, metadata=None, run_id=None)

Convenience wrapper for fail() that always includes exception.

Parameters:

Name Type Description Default
msg str

Human-readable error message

required
code str

Machine-readable error code (must start with ZEO_, ZC_, or the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)

required
exc Exception

Exception that caused the error

required
metadata dict[str, Any] | None

Optional metadata dict

None
run_id str | None

Optional run_id to reuse (should match manifest run_id)

None

Returns:

Type Description
CapabilityResult[T]

CapabilityResult with status=error

Example

try: ... process_video() ... except IOError as e: ... result = CapabilityResult.fail_from_exc( ... msg="Video processing failed", ... code="ZEO_IO_ERROR", ... exc=e ... )

Source code in src/zeo_core/contracts/envelopes/result.py
@classmethod
def fail_from_exc(
    cls,
    msg: str,
    code: str,
    exc: Exception,
    metadata: dict[str, Any] | None = None,
    run_id: str | None = None,
) -> "CapabilityResult[T]":
    """
    Convenience wrapper for fail() that always includes exception.

    Args:
        msg: Human-readable error message
        code: Machine-readable error code (must start with ZEO_, ZC_, or
            the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)
        exc: Exception that caused the error
        metadata: Optional metadata dict
        run_id: Optional run_id to reuse (should match manifest run_id)

    Returns:
        CapabilityResult with status=error

    Example:
        >>> try:
        ...     process_video()
        ... except IOError as e:
        ...     result = CapabilityResult.fail_from_exc(
        ...         msg="Video processing failed",
        ...         code="ZEO_IO_ERROR",
        ...         exc=e
        ...     )
    """
    return cls.fail(
        msg=msg,
        code=code,
        exception=exc,
        metadata=metadata,
        run_id=run_id,
        outcome=CapabilityOutcome.unexpected_exception,
    )

ok classmethod

ok(data, msg='Success', metadata=None, logs=None, duration_sec=None, run_id=None)

Create a successful result.

Parameters:

Name Type Description Default
data T

The result payload

required
msg str

Human-readable success message

'Success'
metadata dict[str, Any] | None

Optional metadata dict

None
logs list[CapabilityLogEvent] | None

Optional log events from execution

None
duration_sec float | None

Execution time in seconds (None if not measured)

None
run_id str | None

Optional run_id to reuse (should match manifest run_id)

None

Returns:

Type Description
CapabilityResult[T]

CapabilityResult with status=success

Example

result = CapabilityResult.ok( ... data={"clips": [...]}, ... msg="Generated 5 clips", ... metadata={"tool": "slice_video", "preset": "fast"} ... )

Source code in src/zeo_core/contracts/envelopes/result.py
@classmethod
def ok(
    cls,
    data: T,
    msg: str = "Success",
    metadata: dict[str, Any] | None = None,
    logs: list[CapabilityLogEvent] | None = None,
    duration_sec: float | None = None,
    run_id: str | None = None,
) -> "CapabilityResult[T]":
    """
    Create a successful result.

    Args:
        data: The result payload
        msg: Human-readable success message
        metadata: Optional metadata dict
        logs: Optional log events from execution
        duration_sec: Execution time in seconds (None if not measured)
        run_id: Optional run_id to reuse (should match manifest run_id)

    Returns:
        CapabilityResult with status=success

    Example:
        >>> result = CapabilityResult.ok(
        ...     data={"clips": [...]},
        ...     msg="Generated 5 clips",
        ...     metadata={"tool": "slice_video", "preset": "fast"}
        ... )
    """
    kwargs = {
        "status": CapabilityStatus.success,
        "outcome": CapabilityOutcome.success,
        "data": data,
        "human_message": msg,
        "metadata": metadata or {},
        "logs": logs or [],
        "duration_sec": duration_sec,
    }
    if run_id is not None:
        kwargs["run_id"] = run_id
    return cls(**kwargs)

skip classmethod

skip(reason, code, metadata=None, run_id=None)

Create a skip result (valid policy decision).

Skips are NOT errors - they represent intentional decisions to skip processing (e.g., video too short, file already exists).

Parameters:

Name Type Description Default
reason str

Human-readable explanation for the skip

required
code str

Machine-readable skip code (must start with ZEO_, ZC_, or the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)

required
metadata dict[str, Any] | None

Optional metadata dict

None
run_id str | None

Optional run_id to reuse (should match manifest run_id)

None

Returns:

Type Description
CapabilityResult[T]

CapabilityResult with status=skipped

Example

result = CapabilityResult.skip( ... reason="Video duration under 10 seconds", ... code="ZEO_VAL_TOO_SHORT" ... )

Source code in src/zeo_core/contracts/envelopes/result.py
@classmethod
def skip(
    cls,
    reason: str,
    code: str,
    metadata: dict[str, Any] | None = None,
    run_id: str | None = None,
) -> "CapabilityResult[T]":
    """
    Create a skip result (valid policy decision).

    Skips are NOT errors - they represent intentional decisions
    to skip processing (e.g., video too short, file already exists).

    Args:
        reason: Human-readable explanation for the skip
        code: Machine-readable skip code (must start with ZEO_, ZC_, or
            the legacy QC_ -- see MACHINE_MESSAGE_PREFIXES)
        metadata: Optional metadata dict
        run_id: Optional run_id to reuse (should match manifest run_id)

    Returns:
        CapabilityResult with status=skipped

    Example:
        >>> result = CapabilityResult.skip(
        ...     reason="Video duration under 10 seconds",
        ...     code="ZEO_VAL_TOO_SHORT"
        ... )
    """
    kwargs = {
        "status": CapabilityStatus.skipped,
        "outcome": CapabilityOutcome.policy_skipped,
        "human_message": reason,
        "machine_message": code,
        "metadata": metadata or {},
    }
    if run_id is not None:
        kwargs["run_id"] = run_id
    return cls(**kwargs)

unavailable classmethod

unavailable(reason, code='ZEO_CAP_UNAVAILABLE', metadata=None, run_id=None)

Create a skipped result because a declared dependency is missing.

Source code in src/zeo_core/contracts/envelopes/result.py
@classmethod
def unavailable(
    cls,
    reason: str,
    code: str = "ZEO_CAP_UNAVAILABLE",
    metadata: dict[str, Any] | None = None,
    run_id: str | None = None,
) -> "CapabilityResult[T]":
    """Create a skipped result because a declared dependency is missing."""
    kwargs: dict[str, Any] = {
        "status": CapabilityStatus.skipped,
        "outcome": CapabilityOutcome.unavailable,
        "human_message": reason,
        "machine_message": code,
        "metadata": metadata or {},
    }
    if run_id is not None:
        kwargs["run_id"] = run_id
    return cls(**kwargs)

validate_machine_message_format classmethod

validate_machine_message_format(v)

Ensure machine_message follows a recognized *_ convention when present.

Source code in src/zeo_core/contracts/envelopes/result.py
@field_validator("machine_message")
@classmethod
def validate_machine_message_format(cls, v: str | None) -> str | None:
    """Ensure machine_message follows a recognized *_ convention when present."""
    if v is not None and not v.startswith(cls.MACHINE_MESSAGE_PREFIXES):
        raise ValueError(
            f"machine_message must start with one of "
            f"{cls.MACHINE_MESSAGE_PREFIXES}, got: {v}. "
            "Use format: ZEO_<AREA>_<DETAIL> (e.g., ZEO_VAL_TOO_SHORT)"
        )
    return v

validate_status_invariants

validate_status_invariants()

Enforce invariants between status and other fields.

This ensures orchestrators can rely on the structure: - Errors always have error objects and machine codes - Successes never have error objects or machine codes - Skips never have error objects but always have machine codes

Source code in src/zeo_core/contracts/envelopes/result.py
@model_validator(mode="after")
def validate_status_invariants(self) -> "CapabilityResult[T]":
    """
    Enforce invariants between status and other fields.

    This ensures orchestrators can rely on the structure:
    - Errors always have error objects and machine codes
    - Successes never have error objects or machine codes
    - Skips never have error objects but always have machine codes
    """
    if self.status == CapabilityStatus.error:
        if self.error is None:
            raise ValueError("status=error requires error field to be present")
        if self.machine_message is None:
            raise ValueError("status=error requires machine_message for branching")

    if self.status == CapabilityStatus.success:
        if self.error is not None:
            raise ValueError("status=success must not have error field")
        if self.machine_message is not None:
            raise ValueError(
                "status=success should not have machine_message "
                "(success is the default path, no special routing needed)"
            )

    if self.status == CapabilityStatus.skipped:
        if self.error is not None:
            raise ValueError(
                "status=skipped must not have error field "
                "(skips are policy decisions, not errors)"
            )
        if self.machine_message is None:
            raise ValueError(
                "status=skipped requires machine_message for branching"
            )

    self._coerce_outcome()
    return self

zeo_core.contracts.EffectKind

Bases: StrEnum

Declared side-effect kind. Declarations are not permission grants.

Source code in src/zeo_core/contracts/common/enums.py
class EffectKind(StrEnum):
    """Declared side-effect kind. Declarations are not permission grants."""

    READ = "read"
    WRITE = "write"
    DELETE = "delete"
    EXTERNAL_COMMUNICATION = "external_communication"
    FINANCIAL = "financial"
    SECURITY_SENSITIVE = "security_sensitive"

Authoring and invocation

zeo_core.tools.capability

capability(*, id, description, effects, examples, error_codes=(), concurrency=ConcurrencyMode.PARALLEL_SAFE, resource_key_fields=(), requirements=None, tags=(), metadata=None, deprecation=None, projection_name=None, guards=(), register_to=None)

Wrap a typed function as a BoundCapability.

Does not register globally unless register_to is an explicit registry.

Source code in src/zeo_core/tools/authoring.py
def capability(
    *,
    id: str,  # noqa: A002 -- canonical keyword matches CapabilityId string form
    description: str,
    effects: Iterable[EffectKind],
    examples: Sequence[CapabilityExample],
    error_codes: Sequence[str] | frozenset[str] = (),
    concurrency: ConcurrencyMode = ConcurrencyMode.PARALLEL_SAFE,
    resource_key_fields: tuple[str, ...] = (),
    requirements: CapabilityRequirements | None = None,
    tags: Sequence[str] | frozenset[str] = (),
    metadata: Mapping[str, JsonValue] | None = None,
    deprecation: CapabilityDeprecation | None = None,
    projection_name: str | None = None,
    guards: Sequence[RequestGuard] = (),
    register_to: object | None = None,
) -> Callable[[Callable[P, R]], Callable[P, R]]:
    """
    Wrap a typed function as a BoundCapability.

    Does not register globally unless ``register_to`` is an explicit registry.
    """

    def decorator(fn: Callable[P, R]) -> Callable[P, R]:
        request_model, response_model = _validate_signature(fn)
        try:
            definition = build_definition(
                capability_id=id,
                description=description,
                request_model=request_model,
                response_model=response_model,
                effects=effects,
                examples=examples,
                error_codes=error_codes,
                concurrency=concurrency,
                resource_key_fields=resource_key_fields,
                requirements=requirements,
                tags=tags,
                metadata=metadata,
                deprecation=deprecation,
                projection_name=projection_name,
            )
        except (ValidationError, ValueError) as exc:
            raise CapabilityAuthoringError(str(exc)) from exc
        bound = BoundCapability(
            definition=definition,
            fn=fn,  # type: ignore[arg-type]
            request_model=request_model,
            guards=guards,
            is_async=inspect.iscoroutinefunction(fn),
        )
        wrapped = functools.wraps(fn)(fn)
        wrapped.__zeo_capability__ = bound  # type: ignore[attr-defined]
        if register_to is not None:
            register = getattr(register_to, "register", None)
            if not callable(register):
                raise CapabilityAuthoringError(
                    "register_to must be a CapabilityRegistry"
                )
            register(bound)
        return wrapped

    return decorator

zeo_core.tools.ToolContext

Bases: BaseModel

Immutable context for tool execution.

The runner constructs this with all required services. Tools receive it as read-only and must not modify it.

METADATA CONTRACT (ENFORCED): - Must be a Mapping (dict-like) - Must be JSON-serializable (validated on construction) - Uses shared normalize_for_json() - same logic as ToolRunner - Primitives: str, int, float, bool, None - Collections: list, dict (string keys only) - Safe types auto-converted: Path→str, datetime→isoformat, Enum→value - Pydantic models: converted via model_dump() (strict isinstance) - Top-level immutable (cannot reassign ctx.metadata) - Nested values not frozen (tools should treat as read-only) - Violations fail immediately with clear error

Source code in src/zeo_core/tools/context.py
class ToolContext(BaseModel):
    """
    Immutable context for tool execution.

    The runner constructs this with all required services.
    Tools receive it as read-only and must not modify it.

    METADATA CONTRACT (ENFORCED):
    - Must be a Mapping (dict-like)
    - Must be JSON-serializable (validated on construction)
    - Uses shared normalize_for_json() - same logic as ToolRunner
    - Primitives: str, int, float, bool, None
    - Collections: list, dict (string keys only)
    - Safe types auto-converted: Path→str, datetime→isoformat, Enum→value
    - Pydantic models: converted via model_dump() (strict isinstance)
    - Top-level immutable (cannot reassign ctx.metadata)
    - Nested values not frozen (tools should treat as read-only)
    - Violations fail immediately with clear error
    """

    # Identity (required)
    run_id: str
    tool_name: str
    tool_version: str

    # Core services (required - runner must provide)
    logger: Any
    fs: Any

    # Directories (required - accepts str | Path, stores as str)
    work_dir: str
    output_dir: str

    # Integration services (optional - top-level immutable via MappingProxyType)
    services: Mapping[str, Any] = Field(default_factory=dict)

    # Metadata (optional - top-level immutable + JSON-safe enforced)
    metadata: Mapping[str, Any] = Field(default_factory=dict)

    # Validators to accept str | Path
    @field_validator("work_dir", "output_dir", mode="before")
    @classmethod
    def normalize_path(cls, v: str | Path) -> str:
        """
        Normalize Path to str for storage.

        Raises:
            TypeError: If value is neither str nor Path
        """
        if isinstance(v, Path):
            return str(v)
        if isinstance(v, str):
            return v
        raise TypeError(
            f"work_dir/output_dir must be str or Path, got {type(v).__name__}"
        )

    # Must-fix B: Guard against None for services
    @field_validator("services", mode="before")
    @classmethod
    def validate_and_normalize_services(
        cls, v: dict[str, Any] | Mapping[str, Any] | None
    ) -> Mapping[str, Any]:
        """
        Validate services is a Mapping and convert to MappingProxyType.

        Must-fix B: Defensive coding - handle None gracefully.
        """
        if v is None:
            return MappingProxyType({})

        if isinstance(v, MappingProxyType):
            return v

        if not isinstance(v, Mapping):
            raise TypeError(
                f"services must be a Mapping (dict-like), got {type(v).__name__}"
            )

        return MappingProxyType(dict(v))

    # Must-fix #2: Strict type checking for metadata
    @field_validator("metadata", mode="before")
    @classmethod
    def validate_and_normalize_metadata(
        cls, v: dict[str, Any] | Mapping[str, Any] | None
    ) -> Mapping[str, Any]:
        """
        Validate metadata is a Mapping and JSON-serializable, then normalize.

        Must-fix #2: Explicitly enforce Mapping type before processing.
        """
        if v is None:
            return MappingProxyType({})

        if isinstance(v, MappingProxyType):
            return v

        # Must-fix #2: Strict type check BEFORE attempting conversion
        if not isinstance(v, Mapping):
            raise TypeError(
                f"metadata must be a Mapping (dict-like), got {type(v).__name__}. "
                f"Example: metadata={{'key': 'value'}}. "
                f"Cannot pass list, string, or other non-mapping types."
            )

        # Now safe to convert to dict
        metadata_dict = dict(v)

        # Use shared normalization logic
        try:
            normalized = normalize_for_json(
                metadata_dict,
                path="metadata",
                allow_pydantic=True,
                allow_string_fallback=False,
                logger=None,
            )
        except TypeError as e:
            raise TypeError(
                f"ToolContext metadata validation failed: {e}. "
                f"Metadata must be JSON-serializable. "
                f"See zeo_core.core.serialization.normalize_for_json for details."
            ) from e

        # Return as immutable
        return MappingProxyType(normalized)

    # Serializers for MappingProxyType
    @field_serializer("services", "metadata")
    def serialize_mapping(self, v: Mapping[str, Any]) -> dict[str, Any]:
        """Convert MappingProxyType to dict for serialization."""
        return dict(v)

    # Configuration: frozen (immutable)
    model_config = ConfigDict(frozen=True, arbitrary_types_allowed=True)

    # Convenience properties for Path usage

    @property
    def work_path(self) -> Path:
        """Get work directory as Path for path computation (no I/O)."""
        return Path(self.work_dir)

    @property
    def output_path(self) -> Path:
        """Get output directory as Path (reference only, no direct writes)."""
        return Path(self.output_dir)

    # Accessor methods (pure - no side effects)

    def require_logger(self) -> Any:  # noqa: ANN401 -- logger is a runner-provided duck-typed service, no fixed protocol exists in this codebase
        """Get logger (guaranteed non-None by runner)."""
        return self.logger

    def require_fs(self) -> Any:  # noqa: ANN401 -- fs service is duck-typed by runner (see mixins/env_init.py MIGRATION COMPAT for the varying .success/.ok, .data/.value contracts it must tolerate)
        """Get filesystem service (guaranteed non-None by runner)."""
        return self.fs

    def get_service(self, name: str) -> Any | None:  # noqa: ANN401 -- services is a heterogeneous registry (Mapping[str, Any]) of arbitrary runner-provided integration instances, looked up dynamically by name
        """Get a service by name (if runner provided it)."""
        return self.services.get(name)

    def require_service(self, name: str) -> Any:  # noqa: ANN401 -- same rationale as get_service: arbitrary integration instance from the heterogeneous services registry
        """
        Get a service by name (raises if missing).

        Raises:
            ValueError: If service not available in context
        """
        service = self.get_service(name)
        if service is None:
            raise ValueError(
                f"Service '{name}' not available in context. "
                f"Runner must provide it in ctx.services. "
                f"Available services: {list(self.services.keys())}"
            )
        return service

    def get_clock(self) -> Any:  # noqa: ANN401 -- optional runner clock; default is constructed by invoke helper
        """Return the injected clock service, if any."""
        return self.get_service("clock")

    def get_cancellation(self) -> Any:  # noqa: ANN401 -- optional cancellation token
        """Return the injected cancellation service, if any."""
        return self.get_service("cancellation")

    def get_artifact_sink(self) -> Any:  # noqa: ANN401 -- optional artifact emitter
        """Return the injected artifact sink, if any."""
        return self.get_service("artifacts")

output_path property

output_path

Get output directory as Path (reference only, no direct writes).

work_path property

work_path

Get work directory as Path for path computation (no I/O).

get_artifact_sink

get_artifact_sink()

Return the injected artifact sink, if any.

Source code in src/zeo_core/tools/context.py
def get_artifact_sink(self) -> Any:  # noqa: ANN401 -- optional artifact emitter
    """Return the injected artifact sink, if any."""
    return self.get_service("artifacts")

get_cancellation

get_cancellation()

Return the injected cancellation service, if any.

Source code in src/zeo_core/tools/context.py
def get_cancellation(self) -> Any:  # noqa: ANN401 -- optional cancellation token
    """Return the injected cancellation service, if any."""
    return self.get_service("cancellation")

get_clock

get_clock()

Return the injected clock service, if any.

Source code in src/zeo_core/tools/context.py
def get_clock(self) -> Any:  # noqa: ANN401 -- optional runner clock; default is constructed by invoke helper
    """Return the injected clock service, if any."""
    return self.get_service("clock")

get_service

get_service(name)

Get a service by name (if runner provided it).

Source code in src/zeo_core/tools/context.py
def get_service(self, name: str) -> Any | None:  # noqa: ANN401 -- services is a heterogeneous registry (Mapping[str, Any]) of arbitrary runner-provided integration instances, looked up dynamically by name
    """Get a service by name (if runner provided it)."""
    return self.services.get(name)

normalize_path classmethod

normalize_path(v)

Normalize Path to str for storage.

Raises:

Type Description
TypeError

If value is neither str nor Path

Source code in src/zeo_core/tools/context.py
@field_validator("work_dir", "output_dir", mode="before")
@classmethod
def normalize_path(cls, v: str | Path) -> str:
    """
    Normalize Path to str for storage.

    Raises:
        TypeError: If value is neither str nor Path
    """
    if isinstance(v, Path):
        return str(v)
    if isinstance(v, str):
        return v
    raise TypeError(
        f"work_dir/output_dir must be str or Path, got {type(v).__name__}"
    )

require_fs

require_fs()

Get filesystem service (guaranteed non-None by runner).

Source code in src/zeo_core/tools/context.py
def require_fs(self) -> Any:  # noqa: ANN401 -- fs service is duck-typed by runner (see mixins/env_init.py MIGRATION COMPAT for the varying .success/.ok, .data/.value contracts it must tolerate)
    """Get filesystem service (guaranteed non-None by runner)."""
    return self.fs

require_logger

require_logger()

Get logger (guaranteed non-None by runner).

Source code in src/zeo_core/tools/context.py
def require_logger(self) -> Any:  # noqa: ANN401 -- logger is a runner-provided duck-typed service, no fixed protocol exists in this codebase
    """Get logger (guaranteed non-None by runner)."""
    return self.logger

require_service

require_service(name)

Get a service by name (raises if missing).

Raises:

Type Description
ValueError

If service not available in context

Source code in src/zeo_core/tools/context.py
def require_service(self, name: str) -> Any:  # noqa: ANN401 -- same rationale as get_service: arbitrary integration instance from the heterogeneous services registry
    """
    Get a service by name (raises if missing).

    Raises:
        ValueError: If service not available in context
    """
    service = self.get_service(name)
    if service is None:
        raise ValueError(
            f"Service '{name}' not available in context. "
            f"Runner must provide it in ctx.services. "
            f"Available services: {list(self.services.keys())}"
        )
    return service

serialize_mapping

serialize_mapping(v)

Convert MappingProxyType to dict for serialization.

Source code in src/zeo_core/tools/context.py
@field_serializer("services", "metadata")
def serialize_mapping(self, v: Mapping[str, Any]) -> dict[str, Any]:
    """Convert MappingProxyType to dict for serialization."""
    return dict(v)

validate_and_normalize_metadata classmethod

validate_and_normalize_metadata(v)

Validate metadata is a Mapping and JSON-serializable, then normalize.

Must-fix #2: Explicitly enforce Mapping type before processing.

Source code in src/zeo_core/tools/context.py
@field_validator("metadata", mode="before")
@classmethod
def validate_and_normalize_metadata(
    cls, v: dict[str, Any] | Mapping[str, Any] | None
) -> Mapping[str, Any]:
    """
    Validate metadata is a Mapping and JSON-serializable, then normalize.

    Must-fix #2: Explicitly enforce Mapping type before processing.
    """
    if v is None:
        return MappingProxyType({})

    if isinstance(v, MappingProxyType):
        return v

    # Must-fix #2: Strict type check BEFORE attempting conversion
    if not isinstance(v, Mapping):
        raise TypeError(
            f"metadata must be a Mapping (dict-like), got {type(v).__name__}. "
            f"Example: metadata={{'key': 'value'}}. "
            f"Cannot pass list, string, or other non-mapping types."
        )

    # Now safe to convert to dict
    metadata_dict = dict(v)

    # Use shared normalization logic
    try:
        normalized = normalize_for_json(
            metadata_dict,
            path="metadata",
            allow_pydantic=True,
            allow_string_fallback=False,
            logger=None,
        )
    except TypeError as e:
        raise TypeError(
            f"ToolContext metadata validation failed: {e}. "
            f"Metadata must be JSON-serializable. "
            f"See zeo_core.core.serialization.normalize_for_json for details."
        ) from e

    # Return as immutable
    return MappingProxyType(normalized)

validate_and_normalize_services classmethod

validate_and_normalize_services(v)

Validate services is a Mapping and convert to MappingProxyType.

Must-fix B: Defensive coding - handle None gracefully.

Source code in src/zeo_core/tools/context.py
@field_validator("services", mode="before")
@classmethod
def validate_and_normalize_services(
    cls, v: dict[str, Any] | Mapping[str, Any] | None
) -> Mapping[str, Any]:
    """
    Validate services is a Mapping and convert to MappingProxyType.

    Must-fix B: Defensive coding - handle None gracefully.
    """
    if v is None:
        return MappingProxyType({})

    if isinstance(v, MappingProxyType):
        return v

    if not isinstance(v, Mapping):
        raise TypeError(
            f"services must be a Mapping (dict-like), got {type(v).__name__}"
        )

    return MappingProxyType(dict(v))

zeo_core.tools.invoke_sync

invoke_sync(capability, request, ctx)
Source code in src/zeo_core/tools/invoke.py
def invoke_sync(
    capability: BoundCapability, request: BaseModel, ctx: ToolContext
) -> CapabilityResult[Any]:
    if context_cancellation(ctx).is_cancelled():
        return CapabilityResult.fail(
            msg="Caller cancellation observed",
            code="ZEO_CAP_CANCELLED",
            outcome=CapabilityOutcome.cancelled,
        )
    if not isinstance(request, capability.request_model):
        try:
            request = capability.request_model.model_validate(
                request if isinstance(request, dict) else request
            )
        except ValidationError as exc:
            return _exception_result(exc)

    guard = _run_guards(capability, request)
    if not guard.ok:
        return CapabilityResult.fail(
            msg=guard.message or "Request rejected by guard",
            code=guard.code or "ZEO_CAP_GUARD_REJECTED",
            outcome=CapabilityOutcome.guard_rejected,
            metadata={"issues": [i.model_dump() for i in guard.issues]},
        )

    if not capability.is_available(ctx):
        missing = missing_requirements(capability.definition, ctx)
        return CapabilityResult.unavailable(
            reason=(
                "Capability unavailable; missing: " + (", ".join(missing) or "unknown")
            ),
        )

    try:
        raw = capability._fn(request, ctx)
        if inspect.isawaitable(raw):
            return CapabilityResult.fail(
                msg="Async capability invoked with invoke_sync",
                code="ZEO_CAP_INVALID_RETURN",
                outcome=CapabilityOutcome.invalid_return,
            )
        return _normalize_return(
            raw, cancelled=context_cancellation(ctx).is_cancelled()
        )
    except BaseException as exc:  # noqa: BLE001 -- convert to structured result
        if context_cancellation(ctx).is_cancelled():
            return CapabilityResult.fail(
                msg="Caller cancellation observed",
                code="ZEO_CAP_CANCELLED",
                exception=exc if isinstance(exc, Exception) else None,
                outcome=CapabilityOutcome.cancelled,
            )
        return _exception_result(exc)

Added in 0.11.0: Runtime and meetings

These APIs ship in 0.11.0. The generic host protocol is a candidate; consult its guide before integration.

zeo_core.contracts.runtime.LaunchContext

Bases: WireModel

Delivered on a private inherited FD, not supplied in request JSON.

Source code in src/zeo_core/contracts/runtime.py
class LaunchContext(WireModel):
    """Delivered on a private inherited FD, not supplied in request JSON."""

    protocol_version: Literal[1]
    provider: ProviderBinding
    attempt: AttemptBinding
    bootstrap_id: Identifier
    admitted_capabilities: tuple[Identifier, ...]
    deadline_unix_ms: Annotated[int, Field(gt=0)]
    # Scope requirements are asserted by Runtime, not discovered from the host env.
    services: tuple[Identifier, ...] = ()
    credentials: tuple[Identifier, ...] = ()
    binaries: tuple[Identifier, ...] = ()
    network_hosts: tuple[Identifier, ...] = ()
    network_allowed: bool = False
    filesystem_read: bool = False
    filesystem_write: bool = False
    workspace: Identifier
    max_message_bytes: Annotated[int, Field(ge=1024, le=1048576)] = 1048576
    rpc_timeout_ms: Annotated[int, Field(ge=1, le=30000)] = 5000
    total_dispatch_budget: Annotated[int, Field(ge=0)] = 0

zeo_core.contracts.runtime.HostResult

Bases: WireModel

Source code in src/zeo_core/contracts/runtime.py
class HostResult(WireModel):
    protocol_version: Literal[1]
    binding: AttemptBinding | None
    state: State
    effect_disposition: Literal["none", "confirmed", "unknown"]
    data: JsonValue = None
    error_code: Identifier | None = None
    # These are Runtime-issued references, never arbitrary local paths.
    artifact_refs: tuple[Identifier, ...] = ()

    @model_validator(mode="after")
    def coherent(self) -> HostResult:
        if self.state == "succeeded" and (
            self.error_code or self.effect_disposition == "unknown"
        ):
            raise ValueError("success cannot contain an error or unknown effect")
        if (
            self.state
            in {
                "failed",
                "refused",
                "protocol_error",
                "unavailable",
                "timed_out",
                "invalid_request",
            }
            and not self.error_code
        ):
            raise ValueError("unsuccessful state requires an error code")
        if self.effect_disposition == "unknown" and self.state not in {
            "needs_reconciliation",
            "protocol_error",
        }:
            raise ValueError("unknown effect requires reconciliation")
        return self

zeo_core.integrations.meetings.MeetingRunner

Runtime owns leases and replay; this runner owns no second execution ledger.

Source code in src/zeo_core/integrations/meetings/runner.py
class MeetingRunner:
    """Runtime owns leases and replay; this runner owns no second execution ledger."""

    def __init__(
        self,
        *,
        runtime: MeetingRuntimePort,
        adapters: AdapterPort,
        clock: Callable[[], datetime] | None = None,
        observe: Callable[[InvocationReceipt, JsonValue], None] | None = None,
    ) -> None:
        self._observe = observe
        self._runtime = runtime
        self._adapters = adapters
        self._clock = clock or (lambda: datetime.now(UTC))

    def execute(
        self, authorization: InvocationAuthorization, payload: dict[str, JsonValue]
    ) -> InvocationReceipt:
        # Freeze arguments before validation and runtime admission.
        frozen = json.loads(json.dumps(payload, allow_nan=False))
        if request_digest(frozen) != authorization.request_sha256:
            raise AdapterRefusalError("meeting request digest mismatch")
        prepared = self._adapters.prepare(authorization, frozen)
        admitted = self._runtime.authorize(authorization)
        if admitted.authorization != authorization:
            raise RuntimeRefusalError("runtime authorization binding mismatch")
        if admitted.disposition == "RECONCILE":
            raise AdapterAmbiguousError("runtime requires reconciliation; no dispatch")
        if admitted.disposition == "REPLAY_FINAL":
            receipt = admitted.receipt
            if receipt is None or not _matches(receipt, authorization):
                raise RuntimeRefusalError("runtime replay receipt mismatch")
            reconciliation = admitted.reconciliation
            if reconciliation is not None:
                if reconciliation.invocation_id != authorization.invocation_id:
                    raise RuntimeRefusalError("runtime reconciliation binding mismatch")
                # A projection of the runtime's two admitted records, not a new receipt.
                return InvocationReceipt.model_validate(
                    {
                        **receipt.model_dump(),
                        "outcome": reconciliation.outcome,
                        "provider_object_id": reconciliation.provider_object_id,
                        "observed_at": reconciliation.observed_at,
                        "reconciled": True,
                    }
                )
            return receipt
        if admitted.receipt is not None or admitted.reconciliation is not None:
            raise RuntimeRefusalError("fresh execution carried a prior result")
        return self._dispatch_and_admit(authorization, prepared)

    def _dispatch_and_admit(
        self, authorization: InvocationAuthorization, prepared: PreparedOperation
    ) -> InvocationReceipt:
        object_id = ""
        response_hash = ""
        response: JsonValue = None
        try:
            object_id, response = prepared()
            response_hash = hashlib.sha256(
                json.dumps(
                    response, sort_keys=True, separators=(",", ":"), allow_nan=False
                ).encode()
            ).hexdigest()
            if authorization.operation.endswith(".create") and not object_id:
                raise AdapterAmbiguousError("provider object identity is missing")
            outcome = "SUCCEEDED"
        except AdapterRefusalError:
            outcome = "REFUSED"
        except Exception:
            outcome = (
                "AMBIGUOUS"
                if authorization.operation
                in {"gmail.draft.create", "notion.page.upsert", "gmail.draft.reconcile"}
                else "REFUSED"
            )
            object_id = ""
            response_hash = ""
        receipt = InvocationReceipt(
            receipt_id="receipt_" + uuid.uuid4().hex,
            invocation_id=authorization.invocation_id,
            operation=authorization.operation,
            resource=authorization.resource,
            request_sha256=authorization.request_sha256,
            outcome=outcome,
            provider_object_id=object_id,
            response_sha256=response_hash,
            observed_at=self._clock(),
        )
        accepted = self._runtime.admit_receipt(receipt)
        if accepted != receipt:
            raise RuntimeRefusalError(
                "runtime did not admit the exact provider receipt"
            )
        if accepted.outcome == "SUCCEEDED" and self._observe is not None:
            self._observe(accepted, response)
        return accepted

Marketing

zeo_core.integrations.hubspot.HubSpotIntegration

Source code in src/zeo_core/integrations/hubspot/service.py
class HubSpotIntegration:
    integration_id = "hubspot.marketing"
    name = "HubSpot Marketing"
    version = "1.0.0"

    def __init__(self, client: HubSpotClient | None = None) -> None:
        self._client = client

    @property
    def client(self) -> HubSpotClient:
        if self._client is None:
            raise ValueError("HubSpot marketing integration is not initialized")
        return self._client

    def initialize(self) -> IntegrationResult[object]:
        if self._client is None:
            token = os.environ.get("HUBSPOT_ACCESS_TOKEN", "")
            if not token.strip():
                return IntegrationResult.error_result(
                    "Set HUBSPOT_ACCESS_TOKEN to a private-app or OAuth access token"
                )
            self._client = HubSpotClient(SecretStr(token))
        return IntegrationResult.success_result(
            message=(
                "HubSpot client configured; live scopes and account entitlement "
                "are unverified"
            )
        )

    def is_available(self) -> bool:
        return self._client is not None

    def close(self) -> None:
        if self._client is not None:
            self._client.close()
            self._client = None

zeo_core.integrations.kit.KitIntegration

Source code in src/zeo_core/integrations/kit/service.py
class KitIntegration:
    integration_id = "kit.marketing"
    name = "Kit Marketing"
    version = "1.0.0"

    def __init__(self, client: KitClient | None = None) -> None:
        self._client = client

    @property
    def client(self) -> KitClient:
        if self._client is None:
            raise ValueError("Kit integration is not initialized")
        return self._client

    def initialize(self) -> IntegrationResult[object]:
        if self._client is None:
            key, token = (
                os.environ.get("KIT_API_KEY", ""),
                os.environ.get("KIT_ACCESS_TOKEN", ""),
            )
            if bool(key.strip()) == bool(token.strip()):
                return IntegrationResult.error_result(
                    "Set exactly one of KIT_API_KEY or KIT_ACCESS_TOKEN"
                )
            self._client = KitClient(
                SecretStr(key) if key.strip() else None,
                access_token=SecretStr(token) if token.strip() else None,
            )
        return IntegrationResult.success_result(
            message=(
                "Kit client configured; live account entitlement and "
                "authentication unverified"
            )
        )

    def is_available(self) -> bool:
        return self._client is not None

    def close(self) -> None:
        if self._client is not None:
            self._client.close()
            self._client = None

Account tracks and execution profiles

zeo_core.integrations.environments.IntegrationEnvironment dataclass

One application process, one mode, explicit provider inputs.

Fixture mode removes all provider variables; applications supply their own controlled transports/services. Neither backend is an operating-system sandbox.

Source code in src/zeo_core/integrations/environments/launcher.py
@dataclass(frozen=True)
class IntegrationEnvironment:
    """One application process, one mode, explicit provider inputs.

    Fixture mode removes all provider variables; applications supply their own
    controlled transports/services. Neither backend is an operating-system sandbox.
    """

    mode: IntegrationMode
    root: Path
    integrations: tuple[str, ...]
    backend: IntegrationBackend = "live"

    def __post_init__(self) -> None:
        if self.mode not in {"test", "production"}:
            raise ValueError("Mode must be test or production")
        if self.backend not in {"live", "fixture"} or (
            self.backend == "fixture" and self.mode != "test"
        ):
            raise ValueError("Fixture execution requires test mode")
        if not self.integrations or len(set(self.integrations)) != len(
            self.integrations
        ):
            raise ValueError("Choose at least one integration without duplicates")
        if any(name not in CATALOG for name in self.integrations):
            raise ValueError("Unknown integration; consult the setup catalog")
        object.__setattr__(self, "root", Path(self.root).expanduser().resolve())

    @property
    def state_dir(self) -> Path:
        if self.backend == "fixture":
            return self.root / "fixtures" / self.mode
        return self.root / self.mode

    @property
    def work_dir(self) -> Path:
        return self.state_dir / "work"

    def prepare(self) -> None:
        """Create separate private directories; refuse aliases into another mode."""
        for path in (
            self.state_dir,
            self.work_dir,
            self.state_dir / "config",
            self.state_dir / "credentials",
            self.state_dir / "tmp",
        ):
            if path.resolve() != path:
                raise ValueError(
                    "Integration environment directories must not be symlinked"
                )
            path.mkdir(parents=True, exist_ok=True, mode=0o700)
        config = self.state_dir / "config" / "integrations.yaml"
        if config.resolve() != config:
            raise ValueError("Integration config must not be symlinked")
        try:
            descriptor = os.open(config, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
        except FileExistsError:
            return
        with os.fdopen(descriptor, "w") as handle:
            handle.write("{}\n")

    def child_environment(
        self, source: Mapping[str, str] | None = None
    ) -> dict[str, str]:
        """Return selected credentials only; never mutate the parent process."""
        source = os.environ if source is None else source
        result = {key: source[key] for key in _PROCESS_VARIABLES if key in source}
        prefix = f"ZEO_{self.mode.upper()}_"
        other = "ZEO_PRODUCTION_" if self.mode == "test" else "ZEO_TEST_"
        allowed = {key for name in self.integrations for key in CATALOG[name].variables}
        prefixes = tuple(
            key for name in self.integrations for key in CATALOG[name].prefixes
        )
        for key, value in source.items():
            if not key.startswith(prefix):
                continue
            target = key.removeprefix(prefix)
            if target not in allowed and not (prefixes and target.startswith(prefixes)):
                # Other selected-mode providers may be configured in the same shell.
                continue
            if self.backend == "fixture":
                continue
            if (
                value
                and target.endswith(_SECRET_SUFFIXES)
                and value == source.get(other + target)
            ):
                raise ValueError(
                    f"Test and production must use different credentials for {target}"
                )
            if (
                value
                and target == "SUPABASE_URL"
                and value == source.get(other + target)
            ):
                raise ValueError(
                    "Test and production must use different SUPABASE_URL values"
                )
            result[target] = value
        result.update(
            {
                "ZEO_INTEGRATION_MODE": self.mode,
                "ZEO_INTEGRATION_BACKEND": self.backend,
                "ZEO_INTEGRATION_STATE_DIR": str(self.state_dir),
                "ZEO_INTEGRATION_IDS": ",".join(self.integrations),
                # Prevent requests from replacing explicit tokens with HOME/.netrc.
                "NETRC": os.devnull,
                "TMPDIR": str(self.state_dir / "tmp"),
                "TMP": str(self.state_dir / "tmp"),
                "TEMP": str(self.state_dir / "tmp"),
            }
        )
        return result

    def run(
        self,
        command: Sequence[str],
        *,
        source: Mapping[str, str] | None = None,
        timeout: float | None = None,
    ) -> int:
        """Run the application without a shell and propagate its exit status."""
        if not command or not command[0]:
            raise ValueError("An application command is required")
        environment = self.child_environment(source)
        self.prepare()
        result = subprocess.run(  # noqa: S603 -- explicit application, no shell
            list(command),
            cwd=self.work_dir,
            env=environment,
            timeout=timeout,
            check=False,
        )
        return result.returncode

child_environment

child_environment(source=None)

Return selected credentials only; never mutate the parent process.

Source code in src/zeo_core/integrations/environments/launcher.py
def child_environment(
    self, source: Mapping[str, str] | None = None
) -> dict[str, str]:
    """Return selected credentials only; never mutate the parent process."""
    source = os.environ if source is None else source
    result = {key: source[key] for key in _PROCESS_VARIABLES if key in source}
    prefix = f"ZEO_{self.mode.upper()}_"
    other = "ZEO_PRODUCTION_" if self.mode == "test" else "ZEO_TEST_"
    allowed = {key for name in self.integrations for key in CATALOG[name].variables}
    prefixes = tuple(
        key for name in self.integrations for key in CATALOG[name].prefixes
    )
    for key, value in source.items():
        if not key.startswith(prefix):
            continue
        target = key.removeprefix(prefix)
        if target not in allowed and not (prefixes and target.startswith(prefixes)):
            # Other selected-mode providers may be configured in the same shell.
            continue
        if self.backend == "fixture":
            continue
        if (
            value
            and target.endswith(_SECRET_SUFFIXES)
            and value == source.get(other + target)
        ):
            raise ValueError(
                f"Test and production must use different credentials for {target}"
            )
        if (
            value
            and target == "SUPABASE_URL"
            and value == source.get(other + target)
        ):
            raise ValueError(
                "Test and production must use different SUPABASE_URL values"
            )
        result[target] = value
    result.update(
        {
            "ZEO_INTEGRATION_MODE": self.mode,
            "ZEO_INTEGRATION_BACKEND": self.backend,
            "ZEO_INTEGRATION_STATE_DIR": str(self.state_dir),
            "ZEO_INTEGRATION_IDS": ",".join(self.integrations),
            # Prevent requests from replacing explicit tokens with HOME/.netrc.
            "NETRC": os.devnull,
            "TMPDIR": str(self.state_dir / "tmp"),
            "TMP": str(self.state_dir / "tmp"),
            "TEMP": str(self.state_dir / "tmp"),
        }
    )
    return result

prepare

prepare()

Create separate private directories; refuse aliases into another mode.

Source code in src/zeo_core/integrations/environments/launcher.py
def prepare(self) -> None:
    """Create separate private directories; refuse aliases into another mode."""
    for path in (
        self.state_dir,
        self.work_dir,
        self.state_dir / "config",
        self.state_dir / "credentials",
        self.state_dir / "tmp",
    ):
        if path.resolve() != path:
            raise ValueError(
                "Integration environment directories must not be symlinked"
            )
        path.mkdir(parents=True, exist_ok=True, mode=0o700)
    config = self.state_dir / "config" / "integrations.yaml"
    if config.resolve() != config:
        raise ValueError("Integration config must not be symlinked")
    try:
        descriptor = os.open(config, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
    except FileExistsError:
        return
    with os.fdopen(descriptor, "w") as handle:
        handle.write("{}\n")

run

run(command, *, source=None, timeout=None)

Run the application without a shell and propagate its exit status.

Source code in src/zeo_core/integrations/environments/launcher.py
def run(
    self,
    command: Sequence[str],
    *,
    source: Mapping[str, str] | None = None,
    timeout: float | None = None,
) -> int:
    """Run the application without a shell and propagate its exit status."""
    if not command or not command[0]:
        raise ValueError("An application command is required")
    environment = self.child_environment(source)
    self.prepare()
    result = subprocess.run(  # noqa: S603 -- explicit application, no shell
        list(command),
        cwd=self.work_dir,
        env=environment,
        timeout=timeout,
        check=False,
    )
    return result.returncode

zeo_core.integrations.hosted.ExecutionProfile

Bases: StrEnum

Explicit execution placements; automatic choice is policy, not a profile.

Source code in src/zeo_core/integrations/hosted/profile.py
class ExecutionProfile(StrEnum):
    """Explicit execution placements; automatic choice is policy, not a profile."""

    FAKE = "fake"
    LOCAL = "local"
    HOSTED = "hosted"
    GOVERNED = "governed"

Gemini requests

zeo_core.integrations.gemini.ImageGenerationRequest

Bases: BaseModel

One billed operation; references and intent are covered by authorization.

Source code in src/zeo_core/integrations/gemini/models.py
class ImageGenerationRequest(BaseModel):
    """One billed operation; references and intent are covered by authorization."""

    model_config = ConfigDict(frozen=True, extra="forbid")

    schema_version: Literal["zeo-image-generation/1"] = "zeo-image-generation/1"
    project_id: str = Field(pattern=IDENTIFIER)
    model: Literal["gemini-3.1-flash-image"] = "gemini-3.1-flash-image"
    prompt: str = Field(min_length=1, max_length=16_000)
    references: tuple[ImageArtifact, ...] = Field(min_length=1, max_length=6)
    output_mime_type: Literal["image/png", "image/jpeg"] = "image/png"
    aspect_ratio: Literal["1:1", "3:2", "2:3", "16:9", "9:16"] = "1:1"
    image_size: Literal["1K", "2K"] = "1K"

    @model_validator(mode="after")
    def _bounded_distinct_references(self) -> ImageGenerationRequest:
        if len({ref.digest for ref in self.references}) != len(self.references):
            raise ValueError("reference images must be distinct")
        if sum(ref.byte_count for ref in self.references) > 32 * 1024 * 1024:
            raise ValueError("reference images exceed the aggregate byte budget")
        if not self.prompt.strip():
            raise ValueError("prompt must contain text")
        return self