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)

Bases: str, Enum

PENDING = 'PENDING'
EXECUTING = 'EXECUTING'
COMPLETED = 'COMPLETED'
class core.durable_delayed_effects.DelayedEffectAction(*values)

Bases: str, Enum

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: object

Decoded durable intent and the result of its latest transition.

Parameters:
action: DelayedEffectAction
state: DelayedEffectState | None
intent_id: str
operation_id: str
effect_type: str
payload_fingerprint: str
payload: dict[str, Any]
due_at_ms: int
created_at_ms: int
owner_token: str
fence: int
attempts: int
receipt: Any = None
error: Any = None
property terminal: bool
property should_execute: bool
exception core.durable_delayed_effects.DelayedEffectConflict

Bases: RuntimeError

An 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: object

Atomic Redis state machine for delayed tool-effect intents.

Parameters:
  • redis (_AsyncRedis)

  • namespace (str)

  • terminal_ttl_seconds (int)

  • execution_lease_seconds (int)

static intent_id_for(operation_id)
Return type:

str

Parameters:

operation_id (str)

key_for_intent(intent_id)
Return type:

str

Parameters:

intent_id (str)

async register(operation_id, effect_type, payload, *, delay_seconds)
Return type:

DelayedEffectIntent

Parameters:
async list_due(*, limit=32)
Return type:

list[str]

Parameters:

limit (int)

async claim(intent_id, owner_token)
Return type:

DelayedEffectIntent

Parameters:
  • intent_id (str)

  • owner_token (str)

async renew(intent)
Return type:

DelayedEffectIntent

Parameters:

intent (DelayedEffectIntent)

async complete(intent, receipt)
Return type:

DelayedEffectIntent

Parameters:
async requeue(intent, *, delay_seconds, error=None, deferred=False)
Return type:

DelayedEffectIntent

Parameters:
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: object

Periodic fenced executor for registered delayed-effect handlers.

Parameters:
  • store (DelayedEffectStore)

  • handlers (Mapping[str, EffectHandler])

  • owner_token (str | None)

  • readiness_check (ReadinessCheck | None)

  • poll_interval_seconds (float)

  • retry_delay_seconds (int)

  • batch_size (int)

async recover_once()

Claim and process one bounded batch of currently-due intents.

Return type:

int

async start()
Return type:

None

async stop()
Return type:

None

core.durable_delayed_effects.delayed_effect_store_for_context(ctx)

Build the scheduler view used by one durable tool handler.

Return type:

DelayedEffectStore

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.

Return type:

DelayedEffectIntent

Parameters: