core.durable_delayed_effects module
Durable scheduler for delayed side effects accepted by tool calls.
Detached asyncio.sleep tasks disappear when a service restarts. This
module replaces them for checkpointed tool executions with a Redis-backed
intent ledger and due-time index. Registration, fenced claims, lease renewal,
retry, and terminal completion are atomic Lua transitions. Intent keys use a
SHA-256 digest of the parent tool operation_id so raw platform identities
never appear in Redis key names.
Only replay-safe handlers belong in this scheduler. A worker whose lease expires may have completed its effect before crashing, so recovery deliberately executes the handler again. Handlers must therefore use set/clear semantics or an external idempotency key and return a JSON-serializable receipt.
- class core.durable_delayed_effects.DelayedEffectState(*values)
-
- PENDING = 'PENDING'
- EXECUTING = 'EXECUTING'
- COMPLETED = 'COMPLETED'
- class core.durable_delayed_effects.DelayedEffectAction(*values)
-
- REGISTERED = 'REGISTERED'
- CACHED = 'CACHED'
- CONFLICT = 'CONFLICT'
- ACQUIRED = 'ACQUIRED'
- RECLAIMED = 'RECLAIMED'
- BUSY = 'BUSY'
- NOT_DUE = 'NOT_DUE'
- RENEWED = 'RENEWED'
- COMPLETED = 'COMPLETED'
- DEFERRED = 'DEFERRED'
- RETRY_SCHEDULED = 'RETRY_SCHEDULED'
- STALE_OWNER = 'STALE_OWNER'
- INVALID_STATE = 'INVALID_STATE'
- MISSING = 'MISSING'
- class core.durable_delayed_effects.DelayedEffectIntent(action, state, intent_id, operation_id, effect_type, payload_fingerprint, payload, due_at_ms, created_at_ms, owner_token, fence, attempts, receipt=None, error=None)
Bases:
objectDecoded durable intent and the result of its latest transition.
- Parameters:
- action: DelayedEffectAction
- state: DelayedEffectState | None
- exception core.durable_delayed_effects.DelayedEffectConflict
Bases:
RuntimeErrorAn operation id was reused for a different delayed effect.
- class core.durable_delayed_effects.DelayedEffectStore(redis, *, namespace='sg:delayed_effect:v1', terminal_ttl_seconds=604800, execution_lease_seconds=300)
Bases:
objectAtomic Redis state machine for delayed tool-effect intents.
- Parameters:
- async register(operation_id, effect_type, payload, *, delay_seconds)
- async claim(intent_id, owner_token)
- Return type:
- Parameters:
- async renew(intent)
- Return type:
- Parameters:
intent (DelayedEffectIntent)
- async complete(intent, receipt)
- Return type:
- Parameters:
intent (DelayedEffectIntent)
receipt (Any)
- async requeue(intent, *, delay_seconds, error=None, deferred=False)
- Return type:
- Parameters:
intent (DelayedEffectIntent)
error (Any)
deferred (bool)
- class core.durable_delayed_effects.DelayedEffectWorker(store, handlers, *, owner_token=None, readiness_check=None, poll_interval_seconds=1.0, retry_delay_seconds=5, batch_size=32)
Bases:
objectPeriodic fenced executor for registered delayed-effect handlers.
- Parameters:
- async recover_once()
Claim and process one bounded batch of currently-due intents.
- Return type:
- core.durable_delayed_effects.delayed_effect_store_for_context(ctx)
Build the scheduler view used by one durable tool handler.
- Return type:
- Parameters:
ctx (Any)
- async core.durable_delayed_effects.register_delayed_effect(ctx, *, effect_type, payload, delay_seconds)
Durably register a delayed effect for the active checkpointed call.