archastro.platform.v1.resources.event_subscriptions

  1# Copyright (c) 2026 ArchAstro Inc. Licensed under the MIT License.
  2# This file is auto-generated by @archastro/sdk-generator. Do not edit.
  3# Content hash: ee8786d0c9ca
  4
  5from __future__ import annotations
  6
  7from typing import Literal, Required, TypedDict
  8
  9from ...runtime.http_client import HttpClient, SyncHttpClient
 10from ...types.common import (
 11    EventSubscription,
 12    EventSubscriptionClaim,
 13    EventSubscriptionHead,
 14    EventSubscriptionPage,
 15    EventSubscriptionQueue,
 16)
 17
 18
 19class EventSubscriptionCreateInput(TypedDict, total=False):
 20    "Create a domain-event subscription"
 21
 22    event_names: Required[list[str]]
 23    "Exact event names to receive."
 24    max_pending_events: int | None
 25    "Queue cap. Defaults to 100; maximum 1000."
 26    name: Required[str]
 27    "Customer-defined subscription name."
 28    retention_seconds: int | None
 29    "Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only."
 30    status: Literal["active", "paused"] | None
 31    "Initial status. Defaults to active."
 32    visibility_timeout_seconds: int | None
 33    "Default claim lease in seconds. Defaults to 300."
 34
 35
 36class EventSubscriptionUpdateInput(TypedDict, total=False):
 37    "Update a domain-event subscription"
 38
 39    event_names: list[str] | None
 40    max_pending_events: int | None
 41    name: str | None
 42    retention_seconds: int | None
 43    "Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only."
 44    status: Literal["active", "paused"] | None
 45    visibility_timeout_seconds: int | None
 46
 47
 48class EventSubscriptionClaimInput(TypedDict, total=False):
 49    "Claim events from a subscription"
 50
 51    consumer_id: str | None
 52    "Consumer identity between 1 and 128 bytes, recorded on the lease for attribution."
 53    max_events: int | None
 54    "Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics."
 55    request_id: str | None
 56    "Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases."
 57    visibility_timeout_seconds: int | None
 58    "Lease duration override between 15 and 3600 seconds."
 59    wait_seconds: int | None
 60    "Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately."
 61
 62
 63class AsyncEventSubscriptionResource:
 64    def __init__(self, http: HttpClient):
 65        self._http = http
 66
 67    async def list(
 68        self, *, page: int | None = None, per_page: int | None = None
 69    ) -> EventSubscriptionPage:
 70        """
 71        List domain-event subscriptions
 72        Lists the subscriptions visible to the caller, with current queue counters.
 73
 74        Args:
 75            page: Page number, starting at 1.
 76            per_page: Subscriptions per page, from 1 through 100.
 77
 78        Returns:
 79            Subscriptions visible to this caller.
 80        """
 81        query: dict[str, object] = {}
 82        if page is not None:
 83            query["page"] = page
 84        if per_page is not None:
 85            query["per_page"] = per_page
 86        return await self._http.request(
 87            "/api/v1/event_subscriptions",
 88            query=query,
 89            response_type=EventSubscriptionPage,
 90        )
 91
 92    async def create(self, input: EventSubscriptionCreateInput) -> EventSubscription:
 93        """
 94        Create a domain-event subscription
 95        Creates a durable subscription that receives exportable domain events matching its exact event names. Matching events fan out into a per-subscription queue bounded by max_pending_events and retention_seconds; consume with claim and acknowledge.
 96
 97        Args:
 98            input: Request body.
 99            input.event_names: Exact event names to receive.
100            input.max_pending_events: Queue cap. Defaults to 100; maximum 1000.
101            input.name: Customer-defined subscription name.
102            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
103            input.status: Initial status. Defaults to active.
104            input.visibility_timeout_seconds: Default claim lease in seconds. Defaults to 300.
105
106        Returns:
107            The new volatile subscription.
108        """
109        return await self._http.request(
110            "/api/v1/event_subscriptions",
111            method="POST",
112            body=input,
113            response_type=EventSubscription,
114        )
115
116    async def delete(self, subscription: str) -> None:
117        """
118        Delete a domain-event subscription
119        Deletes the subscription and every delivery still in its queue.
120
121        Args:
122            subscription: Subscription ID (`esub_...`).
123
124        Returns:
125            No content
126        """
127        await self._http.request(f"/api/v1/event_subscriptions/{subscription}", method="DELETE")
128
129    async def get(self, subscription: str) -> EventSubscription:
130        """
131        Get a domain-event subscription
132        Returns one subscription and its current queue counters.
133
134        Args:
135            subscription: Subscription ID (`esub_...`).
136
137        Returns:
138            Successful response
139        """
140        return await self._http.request(
141            f"/api/v1/event_subscriptions/{subscription}",
142            response_type=EventSubscription,
143        )
144
145    async def update(
146        self, subscription: str, input: EventSubscriptionUpdateInput
147    ) -> EventSubscription:
148        """
149        Update a domain-event subscription
150        Updates matching, status, or queue limits for future fanout. Already queued deliveries remain unless a lower cap trims the oldest entries; retention changes apply to future deliveries only.
151
152        Args:
153            subscription: Subscription ID (`esub_...`).
154            input: Request body.
155            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
156
157        Returns:
158            Successful response
159        """
160        return await self._http.request(
161            f"/api/v1/event_subscriptions/{subscription}",
162            method="PATCH",
163            body=input,
164            response_type=EventSubscription,
165        )
166
167    async def claim(
168        self, subscription: str, input: EventSubscriptionClaimInput
169    ) -> EventSubscriptionClaim:
170        """
171        Claim events from a subscription
172        Atomically leases the oldest unacknowledged delivery. Returns an empty data array when the queue is empty or its head already has an active lease. Passing max_events opts into batch mode, where leased entries are skipped instead of blocking; wait_seconds bounds a long poll on an empty queue.
173
174        Args:
175            subscription: Subscription ID (`esub_...`).
176            input: Request body.
177            input.consumer_id: Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.
178            input.max_events: Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.
179            input.request_id: Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.
180            input.visibility_timeout_seconds: Lease duration override between 15 and 3600 seconds.
181            input.wait_seconds: Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.
182
183        Returns:
184            Successful response
185        """
186        return await self._http.request(
187            f"/api/v1/event_subscriptions/{subscription}/claim",
188            method="POST",
189            body=input,
190            response_type=EventSubscriptionClaim,
191        )
192
193    async def head(self, subscription: str) -> EventSubscriptionHead:
194        """
195        Peek at the head of a subscription queue
196        Returns the oldest unacknowledged delivery without reserving it. This diagnostic read cannot be used as a safe substitute for claim.
197
198        Args:
199            subscription: Subscription ID (`esub_...`).
200
201        Returns:
202            Successful response
203        """
204        return await self._http.request(
205            f"/api/v1/event_subscriptions/{subscription}/head",
206            response_type=EventSubscriptionHead,
207        )
208
209    async def queue(
210        self,
211        subscription: str,
212        *,
213        limit: int | None = None,
214        before_cursor: str | None = None,
215        after_cursor: str | None = None,
216    ) -> EventSubscriptionQueue:
217        """
218        Read a subscription queue
219        Returns a non-reserving, oldest-first view of unacknowledged deliveries, including the queue's loss indicators (queue_epoch, dropped_events_total, dropped_through_position). Use claim to obtain a receipt handle before acknowledging.
220
221        Args:
222            subscription: Subscription ID (`esub_...`).
223            limit: Maximum entries to return, from 1 through 100.
224            before_cursor: Opaque cursor for the preceding page.
225            after_cursor: Opaque cursor for the following page.
226
227        Returns:
228            Successful response
229        """
230        query: dict[str, object] = {}
231        if limit is not None:
232            query["limit"] = limit
233        if before_cursor is not None:
234            query["before_cursor"] = before_cursor
235        if after_cursor is not None:
236            query["after_cursor"] = after_cursor
237        return await self._http.request(
238            f"/api/v1/event_subscriptions/{subscription}/queue",
239            query=query,
240            response_type=EventSubscriptionQueue,
241        )
242
243
244class EventSubscriptionResource:
245    def __init__(self, http: SyncHttpClient):
246        self._http = http
247
248    def list(
249        self, *, page: int | None = None, per_page: int | None = None
250    ) -> EventSubscriptionPage:
251        """
252        List domain-event subscriptions
253        Lists the subscriptions visible to the caller, with current queue counters.
254
255        Args:
256            page: Page number, starting at 1.
257            per_page: Subscriptions per page, from 1 through 100.
258
259        Returns:
260            Subscriptions visible to this caller.
261        """
262        query: dict[str, object] = {}
263        if page is not None:
264            query["page"] = page
265        if per_page is not None:
266            query["per_page"] = per_page
267        return self._http.request(
268            "/api/v1/event_subscriptions",
269            query=query,
270            response_type=EventSubscriptionPage,
271        )
272
273    def create(self, input: EventSubscriptionCreateInput) -> EventSubscription:
274        """
275        Create a domain-event subscription
276        Creates a durable subscription that receives exportable domain events matching its exact event names. Matching events fan out into a per-subscription queue bounded by max_pending_events and retention_seconds; consume with claim and acknowledge.
277
278        Args:
279            input: Request body.
280            input.event_names: Exact event names to receive.
281            input.max_pending_events: Queue cap. Defaults to 100; maximum 1000.
282            input.name: Customer-defined subscription name.
283            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
284            input.status: Initial status. Defaults to active.
285            input.visibility_timeout_seconds: Default claim lease in seconds. Defaults to 300.
286
287        Returns:
288            The new volatile subscription.
289        """
290        return self._http.request(
291            "/api/v1/event_subscriptions",
292            method="POST",
293            body=input,
294            response_type=EventSubscription,
295        )
296
297    def delete(self, subscription: str) -> None:
298        """
299        Delete a domain-event subscription
300        Deletes the subscription and every delivery still in its queue.
301
302        Args:
303            subscription: Subscription ID (`esub_...`).
304
305        Returns:
306            No content
307        """
308        self._http.request(f"/api/v1/event_subscriptions/{subscription}", method="DELETE")
309
310    def get(self, subscription: str) -> EventSubscription:
311        """
312        Get a domain-event subscription
313        Returns one subscription and its current queue counters.
314
315        Args:
316            subscription: Subscription ID (`esub_...`).
317
318        Returns:
319            Successful response
320        """
321        return self._http.request(
322            f"/api/v1/event_subscriptions/{subscription}",
323            response_type=EventSubscription,
324        )
325
326    def update(self, subscription: str, input: EventSubscriptionUpdateInput) -> EventSubscription:
327        """
328        Update a domain-event subscription
329        Updates matching, status, or queue limits for future fanout. Already queued deliveries remain unless a lower cap trims the oldest entries; retention changes apply to future deliveries only.
330
331        Args:
332            subscription: Subscription ID (`esub_...`).
333            input: Request body.
334            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
335
336        Returns:
337            Successful response
338        """
339        return self._http.request(
340            f"/api/v1/event_subscriptions/{subscription}",
341            method="PATCH",
342            body=input,
343            response_type=EventSubscription,
344        )
345
346    def claim(
347        self, subscription: str, input: EventSubscriptionClaimInput
348    ) -> EventSubscriptionClaim:
349        """
350        Claim events from a subscription
351        Atomically leases the oldest unacknowledged delivery. Returns an empty data array when the queue is empty or its head already has an active lease. Passing max_events opts into batch mode, where leased entries are skipped instead of blocking; wait_seconds bounds a long poll on an empty queue.
352
353        Args:
354            subscription: Subscription ID (`esub_...`).
355            input: Request body.
356            input.consumer_id: Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.
357            input.max_events: Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.
358            input.request_id: Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.
359            input.visibility_timeout_seconds: Lease duration override between 15 and 3600 seconds.
360            input.wait_seconds: Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.
361
362        Returns:
363            Successful response
364        """
365        return self._http.request(
366            f"/api/v1/event_subscriptions/{subscription}/claim",
367            method="POST",
368            body=input,
369            response_type=EventSubscriptionClaim,
370        )
371
372    def head(self, subscription: str) -> EventSubscriptionHead:
373        """
374        Peek at the head of a subscription queue
375        Returns the oldest unacknowledged delivery without reserving it. This diagnostic read cannot be used as a safe substitute for claim.
376
377        Args:
378            subscription: Subscription ID (`esub_...`).
379
380        Returns:
381            Successful response
382        """
383        return self._http.request(
384            f"/api/v1/event_subscriptions/{subscription}/head",
385            response_type=EventSubscriptionHead,
386        )
387
388    def queue(
389        self,
390        subscription: str,
391        *,
392        limit: int | None = None,
393        before_cursor: str | None = None,
394        after_cursor: str | None = None,
395    ) -> EventSubscriptionQueue:
396        """
397        Read a subscription queue
398        Returns a non-reserving, oldest-first view of unacknowledged deliveries, including the queue's loss indicators (queue_epoch, dropped_events_total, dropped_through_position). Use claim to obtain a receipt handle before acknowledging.
399
400        Args:
401            subscription: Subscription ID (`esub_...`).
402            limit: Maximum entries to return, from 1 through 100.
403            before_cursor: Opaque cursor for the preceding page.
404            after_cursor: Opaque cursor for the following page.
405
406        Returns:
407            Successful response
408        """
409        query: dict[str, object] = {}
410        if limit is not None:
411            query["limit"] = limit
412        if before_cursor is not None:
413            query["before_cursor"] = before_cursor
414        if after_cursor is not None:
415            query["after_cursor"] = after_cursor
416        return self._http.request(
417            f"/api/v1/event_subscriptions/{subscription}/queue",
418            query=query,
419            response_type=EventSubscriptionQueue,
420        )
class EventSubscriptionCreateInput(typing.TypedDict):
20class EventSubscriptionCreateInput(TypedDict, total=False):
21    "Create a domain-event subscription"
22
23    event_names: Required[list[str]]
24    "Exact event names to receive."
25    max_pending_events: int | None
26    "Queue cap. Defaults to 100; maximum 1000."
27    name: Required[str]
28    "Customer-defined subscription name."
29    retention_seconds: int | None
30    "Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only."
31    status: Literal["active", "paused"] | None
32    "Initial status. Defaults to active."
33    visibility_timeout_seconds: int | None
34    "Default claim lease in seconds. Defaults to 300."

Create a domain-event subscription

event_names: Required[list[str]]

Exact event names to receive.

max_pending_events: int | None

Queue cap. Defaults to 100; maximum 1000.

name: Required[str]

Customer-defined subscription name.

retention_seconds: int | None

Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.

status: Optional[Literal['active', 'paused']]

Initial status. Defaults to active.

visibility_timeout_seconds: int | None

Default claim lease in seconds. Defaults to 300.

class EventSubscriptionUpdateInput(typing.TypedDict):
37class EventSubscriptionUpdateInput(TypedDict, total=False):
38    "Update a domain-event subscription"
39
40    event_names: list[str] | None
41    max_pending_events: int | None
42    name: str | None
43    retention_seconds: int | None
44    "Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only."
45    status: Literal["active", "paused"] | None
46    visibility_timeout_seconds: int | None

Update a domain-event subscription

event_names: list[str] | None
max_pending_events: int | None
name: str | None
retention_seconds: int | None

Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.

status: Optional[Literal['active', 'paused']]
visibility_timeout_seconds: int | None
class EventSubscriptionClaimInput(typing.TypedDict):
49class EventSubscriptionClaimInput(TypedDict, total=False):
50    "Claim events from a subscription"
51
52    consumer_id: str | None
53    "Consumer identity between 1 and 128 bytes, recorded on the lease for attribution."
54    max_events: int | None
55    "Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics."
56    request_id: str | None
57    "Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases."
58    visibility_timeout_seconds: int | None
59    "Lease duration override between 15 and 3600 seconds."
60    wait_seconds: int | None
61    "Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately."

Claim events from a subscription

consumer_id: str | None

Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.

max_events: int | None

Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.

request_id: str | None

Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.

visibility_timeout_seconds: int | None

Lease duration override between 15 and 3600 seconds.

wait_seconds: int | None

Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.

class AsyncEventSubscriptionResource:
 64class AsyncEventSubscriptionResource:
 65    def __init__(self, http: HttpClient):
 66        self._http = http
 67
 68    async def list(
 69        self, *, page: int | None = None, per_page: int | None = None
 70    ) -> EventSubscriptionPage:
 71        """
 72        List domain-event subscriptions
 73        Lists the subscriptions visible to the caller, with current queue counters.
 74
 75        Args:
 76            page: Page number, starting at 1.
 77            per_page: Subscriptions per page, from 1 through 100.
 78
 79        Returns:
 80            Subscriptions visible to this caller.
 81        """
 82        query: dict[str, object] = {}
 83        if page is not None:
 84            query["page"] = page
 85        if per_page is not None:
 86            query["per_page"] = per_page
 87        return await self._http.request(
 88            "/api/v1/event_subscriptions",
 89            query=query,
 90            response_type=EventSubscriptionPage,
 91        )
 92
 93    async def create(self, input: EventSubscriptionCreateInput) -> EventSubscription:
 94        """
 95        Create a domain-event subscription
 96        Creates a durable subscription that receives exportable domain events matching its exact event names. Matching events fan out into a per-subscription queue bounded by max_pending_events and retention_seconds; consume with claim and acknowledge.
 97
 98        Args:
 99            input: Request body.
100            input.event_names: Exact event names to receive.
101            input.max_pending_events: Queue cap. Defaults to 100; maximum 1000.
102            input.name: Customer-defined subscription name.
103            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
104            input.status: Initial status. Defaults to active.
105            input.visibility_timeout_seconds: Default claim lease in seconds. Defaults to 300.
106
107        Returns:
108            The new volatile subscription.
109        """
110        return await self._http.request(
111            "/api/v1/event_subscriptions",
112            method="POST",
113            body=input,
114            response_type=EventSubscription,
115        )
116
117    async def delete(self, subscription: str) -> None:
118        """
119        Delete a domain-event subscription
120        Deletes the subscription and every delivery still in its queue.
121
122        Args:
123            subscription: Subscription ID (`esub_...`).
124
125        Returns:
126            No content
127        """
128        await self._http.request(f"/api/v1/event_subscriptions/{subscription}", method="DELETE")
129
130    async def get(self, subscription: str) -> EventSubscription:
131        """
132        Get a domain-event subscription
133        Returns one subscription and its current queue counters.
134
135        Args:
136            subscription: Subscription ID (`esub_...`).
137
138        Returns:
139            Successful response
140        """
141        return await self._http.request(
142            f"/api/v1/event_subscriptions/{subscription}",
143            response_type=EventSubscription,
144        )
145
146    async def update(
147        self, subscription: str, input: EventSubscriptionUpdateInput
148    ) -> EventSubscription:
149        """
150        Update a domain-event subscription
151        Updates matching, status, or queue limits for future fanout. Already queued deliveries remain unless a lower cap trims the oldest entries; retention changes apply to future deliveries only.
152
153        Args:
154            subscription: Subscription ID (`esub_...`).
155            input: Request body.
156            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
157
158        Returns:
159            Successful response
160        """
161        return await self._http.request(
162            f"/api/v1/event_subscriptions/{subscription}",
163            method="PATCH",
164            body=input,
165            response_type=EventSubscription,
166        )
167
168    async def claim(
169        self, subscription: str, input: EventSubscriptionClaimInput
170    ) -> EventSubscriptionClaim:
171        """
172        Claim events from a subscription
173        Atomically leases the oldest unacknowledged delivery. Returns an empty data array when the queue is empty or its head already has an active lease. Passing max_events opts into batch mode, where leased entries are skipped instead of blocking; wait_seconds bounds a long poll on an empty queue.
174
175        Args:
176            subscription: Subscription ID (`esub_...`).
177            input: Request body.
178            input.consumer_id: Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.
179            input.max_events: Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.
180            input.request_id: Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.
181            input.visibility_timeout_seconds: Lease duration override between 15 and 3600 seconds.
182            input.wait_seconds: Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.
183
184        Returns:
185            Successful response
186        """
187        return await self._http.request(
188            f"/api/v1/event_subscriptions/{subscription}/claim",
189            method="POST",
190            body=input,
191            response_type=EventSubscriptionClaim,
192        )
193
194    async def head(self, subscription: str) -> EventSubscriptionHead:
195        """
196        Peek at the head of a subscription queue
197        Returns the oldest unacknowledged delivery without reserving it. This diagnostic read cannot be used as a safe substitute for claim.
198
199        Args:
200            subscription: Subscription ID (`esub_...`).
201
202        Returns:
203            Successful response
204        """
205        return await self._http.request(
206            f"/api/v1/event_subscriptions/{subscription}/head",
207            response_type=EventSubscriptionHead,
208        )
209
210    async def queue(
211        self,
212        subscription: str,
213        *,
214        limit: int | None = None,
215        before_cursor: str | None = None,
216        after_cursor: str | None = None,
217    ) -> EventSubscriptionQueue:
218        """
219        Read a subscription queue
220        Returns a non-reserving, oldest-first view of unacknowledged deliveries, including the queue's loss indicators (queue_epoch, dropped_events_total, dropped_through_position). Use claim to obtain a receipt handle before acknowledging.
221
222        Args:
223            subscription: Subscription ID (`esub_...`).
224            limit: Maximum entries to return, from 1 through 100.
225            before_cursor: Opaque cursor for the preceding page.
226            after_cursor: Opaque cursor for the following page.
227
228        Returns:
229            Successful response
230        """
231        query: dict[str, object] = {}
232        if limit is not None:
233            query["limit"] = limit
234        if before_cursor is not None:
235            query["before_cursor"] = before_cursor
236        if after_cursor is not None:
237            query["after_cursor"] = after_cursor
238        return await self._http.request(
239            f"/api/v1/event_subscriptions/{subscription}/queue",
240            query=query,
241            response_type=EventSubscriptionQueue,
242        )
AsyncEventSubscriptionResource(http: archastro.platform.runtime.http_client.HttpClient)
65    def __init__(self, http: HttpClient):
66        self._http = http
async def list( self, *, page: int | None = None, per_page: int | None = None) -> archastro.platform.types.common.EventSubscriptionPage:
68    async def list(
69        self, *, page: int | None = None, per_page: int | None = None
70    ) -> EventSubscriptionPage:
71        """
72        List domain-event subscriptions
73        Lists the subscriptions visible to the caller, with current queue counters.
74
75        Args:
76            page: Page number, starting at 1.
77            per_page: Subscriptions per page, from 1 through 100.
78
79        Returns:
80            Subscriptions visible to this caller.
81        """
82        query: dict[str, object] = {}
83        if page is not None:
84            query["page"] = page
85        if per_page is not None:
86            query["per_page"] = per_page
87        return await self._http.request(
88            "/api/v1/event_subscriptions",
89            query=query,
90            response_type=EventSubscriptionPage,
91        )

List domain-event subscriptions Lists the subscriptions visible to the caller, with current queue counters.

Arguments:
  • page: Page number, starting at 1.
  • per_page: Subscriptions per page, from 1 through 100.
Returns:

Subscriptions visible to this caller.

 93    async def create(self, input: EventSubscriptionCreateInput) -> EventSubscription:
 94        """
 95        Create a domain-event subscription
 96        Creates a durable subscription that receives exportable domain events matching its exact event names. Matching events fan out into a per-subscription queue bounded by max_pending_events and retention_seconds; consume with claim and acknowledge.
 97
 98        Args:
 99            input: Request body.
100            input.event_names: Exact event names to receive.
101            input.max_pending_events: Queue cap. Defaults to 100; maximum 1000.
102            input.name: Customer-defined subscription name.
103            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
104            input.status: Initial status. Defaults to active.
105            input.visibility_timeout_seconds: Default claim lease in seconds. Defaults to 300.
106
107        Returns:
108            The new volatile subscription.
109        """
110        return await self._http.request(
111            "/api/v1/event_subscriptions",
112            method="POST",
113            body=input,
114            response_type=EventSubscription,
115        )

Create a domain-event subscription Creates a durable subscription that receives exportable domain events matching its exact event names. Matching events fan out into a per-subscription queue bounded by max_pending_events and retention_seconds; consume with claim and acknowledge.

Arguments:
  • input: Request body.
  • input.event_names: Exact event names to receive.
  • input.max_pending_events: Queue cap. Defaults to 100; maximum 1000.
  • input.name: Customer-defined subscription name.
  • input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
  • input.status: Initial status. Defaults to active.
  • input.visibility_timeout_seconds: Default claim lease in seconds. Defaults to 300.
Returns:

The new volatile subscription.

async def delete(self, subscription: str) -> None:
117    async def delete(self, subscription: str) -> None:
118        """
119        Delete a domain-event subscription
120        Deletes the subscription and every delivery still in its queue.
121
122        Args:
123            subscription: Subscription ID (`esub_...`).
124
125        Returns:
126            No content
127        """
128        await self._http.request(f"/api/v1/event_subscriptions/{subscription}", method="DELETE")

Delete a domain-event subscription Deletes the subscription and every delivery still in its queue.

Arguments:
  • subscription: Subscription ID (esub_...).
Returns:

No content

async def get( self, subscription: str) -> archastro.platform.types.common.EventSubscription:
130    async def get(self, subscription: str) -> EventSubscription:
131        """
132        Get a domain-event subscription
133        Returns one subscription and its current queue counters.
134
135        Args:
136            subscription: Subscription ID (`esub_...`).
137
138        Returns:
139            Successful response
140        """
141        return await self._http.request(
142            f"/api/v1/event_subscriptions/{subscription}",
143            response_type=EventSubscription,
144        )

Get a domain-event subscription Returns one subscription and its current queue counters.

Arguments:
  • subscription: Subscription ID (esub_...).
Returns:

Successful response

async def update( self, subscription: str, input: EventSubscriptionUpdateInput) -> archastro.platform.types.common.EventSubscription:
146    async def update(
147        self, subscription: str, input: EventSubscriptionUpdateInput
148    ) -> EventSubscription:
149        """
150        Update a domain-event subscription
151        Updates matching, status, or queue limits for future fanout. Already queued deliveries remain unless a lower cap trims the oldest entries; retention changes apply to future deliveries only.
152
153        Args:
154            subscription: Subscription ID (`esub_...`).
155            input: Request body.
156            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
157
158        Returns:
159            Successful response
160        """
161        return await self._http.request(
162            f"/api/v1/event_subscriptions/{subscription}",
163            method="PATCH",
164            body=input,
165            response_type=EventSubscription,
166        )

Update a domain-event subscription Updates matching, status, or queue limits for future fanout. Already queued deliveries remain unless a lower cap trims the oldest entries; retention changes apply to future deliveries only.

Arguments:
  • subscription: Subscription ID (esub_...).
  • input: Request body.
  • input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
Returns:

Successful response

async def claim( self, subscription: str, input: EventSubscriptionClaimInput) -> archastro.platform.types.common.EventSubscriptionClaim:
168    async def claim(
169        self, subscription: str, input: EventSubscriptionClaimInput
170    ) -> EventSubscriptionClaim:
171        """
172        Claim events from a subscription
173        Atomically leases the oldest unacknowledged delivery. Returns an empty data array when the queue is empty or its head already has an active lease. Passing max_events opts into batch mode, where leased entries are skipped instead of blocking; wait_seconds bounds a long poll on an empty queue.
174
175        Args:
176            subscription: Subscription ID (`esub_...`).
177            input: Request body.
178            input.consumer_id: Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.
179            input.max_events: Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.
180            input.request_id: Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.
181            input.visibility_timeout_seconds: Lease duration override between 15 and 3600 seconds.
182            input.wait_seconds: Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.
183
184        Returns:
185            Successful response
186        """
187        return await self._http.request(
188            f"/api/v1/event_subscriptions/{subscription}/claim",
189            method="POST",
190            body=input,
191            response_type=EventSubscriptionClaim,
192        )

Claim events from a subscription Atomically leases the oldest unacknowledged delivery. Returns an empty data array when the queue is empty or its head already has an active lease. Passing max_events opts into batch mode, where leased entries are skipped instead of blocking; wait_seconds bounds a long poll on an empty queue.

Arguments:
  • subscription: Subscription ID (esub_...).
  • input: Request body.
  • input.consumer_id: Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.
  • input.max_events: Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.
  • input.request_id: Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.
  • input.visibility_timeout_seconds: Lease duration override between 15 and 3600 seconds.
  • input.wait_seconds: Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.
Returns:

Successful response

async def head( self, subscription: str) -> archastro.platform.types.common.EventSubscriptionHead:
194    async def head(self, subscription: str) -> EventSubscriptionHead:
195        """
196        Peek at the head of a subscription queue
197        Returns the oldest unacknowledged delivery without reserving it. This diagnostic read cannot be used as a safe substitute for claim.
198
199        Args:
200            subscription: Subscription ID (`esub_...`).
201
202        Returns:
203            Successful response
204        """
205        return await self._http.request(
206            f"/api/v1/event_subscriptions/{subscription}/head",
207            response_type=EventSubscriptionHead,
208        )

Peek at the head of a subscription queue Returns the oldest unacknowledged delivery without reserving it. This diagnostic read cannot be used as a safe substitute for claim.

Arguments:
  • subscription: Subscription ID (esub_...).
Returns:

Successful response

async def queue( self, subscription: str, *, limit: int | None = None, before_cursor: str | None = None, after_cursor: str | None = None) -> archastro.platform.types.common.EventSubscriptionQueue:
210    async def queue(
211        self,
212        subscription: str,
213        *,
214        limit: int | None = None,
215        before_cursor: str | None = None,
216        after_cursor: str | None = None,
217    ) -> EventSubscriptionQueue:
218        """
219        Read a subscription queue
220        Returns a non-reserving, oldest-first view of unacknowledged deliveries, including the queue's loss indicators (queue_epoch, dropped_events_total, dropped_through_position). Use claim to obtain a receipt handle before acknowledging.
221
222        Args:
223            subscription: Subscription ID (`esub_...`).
224            limit: Maximum entries to return, from 1 through 100.
225            before_cursor: Opaque cursor for the preceding page.
226            after_cursor: Opaque cursor for the following page.
227
228        Returns:
229            Successful response
230        """
231        query: dict[str, object] = {}
232        if limit is not None:
233            query["limit"] = limit
234        if before_cursor is not None:
235            query["before_cursor"] = before_cursor
236        if after_cursor is not None:
237            query["after_cursor"] = after_cursor
238        return await self._http.request(
239            f"/api/v1/event_subscriptions/{subscription}/queue",
240            query=query,
241            response_type=EventSubscriptionQueue,
242        )

Read a subscription queue Returns a non-reserving, oldest-first view of unacknowledged deliveries, including the queue's loss indicators (queue_epoch, dropped_events_total, dropped_through_position). Use claim to obtain a receipt handle before acknowledging.

Arguments:
  • subscription: Subscription ID (esub_...).
  • limit: Maximum entries to return, from 1 through 100.
  • before_cursor: Opaque cursor for the preceding page.
  • after_cursor: Opaque cursor for the following page.
Returns:

Successful response

class EventSubscriptionResource:
245class EventSubscriptionResource:
246    def __init__(self, http: SyncHttpClient):
247        self._http = http
248
249    def list(
250        self, *, page: int | None = None, per_page: int | None = None
251    ) -> EventSubscriptionPage:
252        """
253        List domain-event subscriptions
254        Lists the subscriptions visible to the caller, with current queue counters.
255
256        Args:
257            page: Page number, starting at 1.
258            per_page: Subscriptions per page, from 1 through 100.
259
260        Returns:
261            Subscriptions visible to this caller.
262        """
263        query: dict[str, object] = {}
264        if page is not None:
265            query["page"] = page
266        if per_page is not None:
267            query["per_page"] = per_page
268        return self._http.request(
269            "/api/v1/event_subscriptions",
270            query=query,
271            response_type=EventSubscriptionPage,
272        )
273
274    def create(self, input: EventSubscriptionCreateInput) -> EventSubscription:
275        """
276        Create a domain-event subscription
277        Creates a durable subscription that receives exportable domain events matching its exact event names. Matching events fan out into a per-subscription queue bounded by max_pending_events and retention_seconds; consume with claim and acknowledge.
278
279        Args:
280            input: Request body.
281            input.event_names: Exact event names to receive.
282            input.max_pending_events: Queue cap. Defaults to 100; maximum 1000.
283            input.name: Customer-defined subscription name.
284            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
285            input.status: Initial status. Defaults to active.
286            input.visibility_timeout_seconds: Default claim lease in seconds. Defaults to 300.
287
288        Returns:
289            The new volatile subscription.
290        """
291        return self._http.request(
292            "/api/v1/event_subscriptions",
293            method="POST",
294            body=input,
295            response_type=EventSubscription,
296        )
297
298    def delete(self, subscription: str) -> None:
299        """
300        Delete a domain-event subscription
301        Deletes the subscription and every delivery still in its queue.
302
303        Args:
304            subscription: Subscription ID (`esub_...`).
305
306        Returns:
307            No content
308        """
309        self._http.request(f"/api/v1/event_subscriptions/{subscription}", method="DELETE")
310
311    def get(self, subscription: str) -> EventSubscription:
312        """
313        Get a domain-event subscription
314        Returns one subscription and its current queue counters.
315
316        Args:
317            subscription: Subscription ID (`esub_...`).
318
319        Returns:
320            Successful response
321        """
322        return self._http.request(
323            f"/api/v1/event_subscriptions/{subscription}",
324            response_type=EventSubscription,
325        )
326
327    def update(self, subscription: str, input: EventSubscriptionUpdateInput) -> EventSubscription:
328        """
329        Update a domain-event subscription
330        Updates matching, status, or queue limits for future fanout. Already queued deliveries remain unless a lower cap trims the oldest entries; retention changes apply to future deliveries only.
331
332        Args:
333            subscription: Subscription ID (`esub_...`).
334            input: Request body.
335            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
336
337        Returns:
338            Successful response
339        """
340        return self._http.request(
341            f"/api/v1/event_subscriptions/{subscription}",
342            method="PATCH",
343            body=input,
344            response_type=EventSubscription,
345        )
346
347    def claim(
348        self, subscription: str, input: EventSubscriptionClaimInput
349    ) -> EventSubscriptionClaim:
350        """
351        Claim events from a subscription
352        Atomically leases the oldest unacknowledged delivery. Returns an empty data array when the queue is empty or its head already has an active lease. Passing max_events opts into batch mode, where leased entries are skipped instead of blocking; wait_seconds bounds a long poll on an empty queue.
353
354        Args:
355            subscription: Subscription ID (`esub_...`).
356            input: Request body.
357            input.consumer_id: Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.
358            input.max_events: Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.
359            input.request_id: Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.
360            input.visibility_timeout_seconds: Lease duration override between 15 and 3600 seconds.
361            input.wait_seconds: Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.
362
363        Returns:
364            Successful response
365        """
366        return self._http.request(
367            f"/api/v1/event_subscriptions/{subscription}/claim",
368            method="POST",
369            body=input,
370            response_type=EventSubscriptionClaim,
371        )
372
373    def head(self, subscription: str) -> EventSubscriptionHead:
374        """
375        Peek at the head of a subscription queue
376        Returns the oldest unacknowledged delivery without reserving it. This diagnostic read cannot be used as a safe substitute for claim.
377
378        Args:
379            subscription: Subscription ID (`esub_...`).
380
381        Returns:
382            Successful response
383        """
384        return self._http.request(
385            f"/api/v1/event_subscriptions/{subscription}/head",
386            response_type=EventSubscriptionHead,
387        )
388
389    def queue(
390        self,
391        subscription: str,
392        *,
393        limit: int | None = None,
394        before_cursor: str | None = None,
395        after_cursor: str | None = None,
396    ) -> EventSubscriptionQueue:
397        """
398        Read a subscription queue
399        Returns a non-reserving, oldest-first view of unacknowledged deliveries, including the queue's loss indicators (queue_epoch, dropped_events_total, dropped_through_position). Use claim to obtain a receipt handle before acknowledging.
400
401        Args:
402            subscription: Subscription ID (`esub_...`).
403            limit: Maximum entries to return, from 1 through 100.
404            before_cursor: Opaque cursor for the preceding page.
405            after_cursor: Opaque cursor for the following page.
406
407        Returns:
408            Successful response
409        """
410        query: dict[str, object] = {}
411        if limit is not None:
412            query["limit"] = limit
413        if before_cursor is not None:
414            query["before_cursor"] = before_cursor
415        if after_cursor is not None:
416            query["after_cursor"] = after_cursor
417        return self._http.request(
418            f"/api/v1/event_subscriptions/{subscription}/queue",
419            query=query,
420            response_type=EventSubscriptionQueue,
421        )
EventSubscriptionResource(http: archastro.platform.runtime.http_client.SyncHttpClient)
246    def __init__(self, http: SyncHttpClient):
247        self._http = http
def list( self, *, page: int | None = None, per_page: int | None = None) -> archastro.platform.types.common.EventSubscriptionPage:
249    def list(
250        self, *, page: int | None = None, per_page: int | None = None
251    ) -> EventSubscriptionPage:
252        """
253        List domain-event subscriptions
254        Lists the subscriptions visible to the caller, with current queue counters.
255
256        Args:
257            page: Page number, starting at 1.
258            per_page: Subscriptions per page, from 1 through 100.
259
260        Returns:
261            Subscriptions visible to this caller.
262        """
263        query: dict[str, object] = {}
264        if page is not None:
265            query["page"] = page
266        if per_page is not None:
267            query["per_page"] = per_page
268        return self._http.request(
269            "/api/v1/event_subscriptions",
270            query=query,
271            response_type=EventSubscriptionPage,
272        )

List domain-event subscriptions Lists the subscriptions visible to the caller, with current queue counters.

Arguments:
  • page: Page number, starting at 1.
  • per_page: Subscriptions per page, from 1 through 100.
Returns:

Subscriptions visible to this caller.

274    def create(self, input: EventSubscriptionCreateInput) -> EventSubscription:
275        """
276        Create a domain-event subscription
277        Creates a durable subscription that receives exportable domain events matching its exact event names. Matching events fan out into a per-subscription queue bounded by max_pending_events and retention_seconds; consume with claim and acknowledge.
278
279        Args:
280            input: Request body.
281            input.event_names: Exact event names to receive.
282            input.max_pending_events: Queue cap. Defaults to 100; maximum 1000.
283            input.name: Customer-defined subscription name.
284            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
285            input.status: Initial status. Defaults to active.
286            input.visibility_timeout_seconds: Default claim lease in seconds. Defaults to 300.
287
288        Returns:
289            The new volatile subscription.
290        """
291        return self._http.request(
292            "/api/v1/event_subscriptions",
293            method="POST",
294            body=input,
295            response_type=EventSubscription,
296        )

Create a domain-event subscription Creates a durable subscription that receives exportable domain events matching its exact event names. Matching events fan out into a per-subscription queue bounded by max_pending_events and retention_seconds; consume with claim and acknowledge.

Arguments:
  • input: Request body.
  • input.event_names: Exact event names to receive.
  • input.max_pending_events: Queue cap. Defaults to 100; maximum 1000.
  • input.name: Customer-defined subscription name.
  • input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
  • input.status: Initial status. Defaults to active.
  • input.visibility_timeout_seconds: Default claim lease in seconds. Defaults to 300.
Returns:

The new volatile subscription.

def delete(self, subscription: str) -> None:
298    def delete(self, subscription: str) -> None:
299        """
300        Delete a domain-event subscription
301        Deletes the subscription and every delivery still in its queue.
302
303        Args:
304            subscription: Subscription ID (`esub_...`).
305
306        Returns:
307            No content
308        """
309        self._http.request(f"/api/v1/event_subscriptions/{subscription}", method="DELETE")

Delete a domain-event subscription Deletes the subscription and every delivery still in its queue.

Arguments:
  • subscription: Subscription ID (esub_...).
Returns:

No content

def get( self, subscription: str) -> archastro.platform.types.common.EventSubscription:
311    def get(self, subscription: str) -> EventSubscription:
312        """
313        Get a domain-event subscription
314        Returns one subscription and its current queue counters.
315
316        Args:
317            subscription: Subscription ID (`esub_...`).
318
319        Returns:
320            Successful response
321        """
322        return self._http.request(
323            f"/api/v1/event_subscriptions/{subscription}",
324            response_type=EventSubscription,
325        )

Get a domain-event subscription Returns one subscription and its current queue counters.

Arguments:
  • subscription: Subscription ID (esub_...).
Returns:

Successful response

def update( self, subscription: str, input: EventSubscriptionUpdateInput) -> archastro.platform.types.common.EventSubscription:
327    def update(self, subscription: str, input: EventSubscriptionUpdateInput) -> EventSubscription:
328        """
329        Update a domain-event subscription
330        Updates matching, status, or queue limits for future fanout. Already queued deliveries remain unless a lower cap trims the oldest entries; retention changes apply to future deliveries only.
331
332        Args:
333            subscription: Subscription ID (`esub_...`).
334            input: Request body.
335            input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
336
337        Returns:
338            Successful response
339        """
340        return self._http.request(
341            f"/api/v1/event_subscriptions/{subscription}",
342            method="PATCH",
343            body=input,
344            response_type=EventSubscription,
345        )

Update a domain-event subscription Updates matching, status, or queue limits for future fanout. Already queued deliveries remain unless a lower cap trims the oldest entries; retention changes apply to future deliveries only.

Arguments:
  • subscription: Subscription ID (esub_...).
  • input: Request body.
  • input.retention_seconds: Delivery retention in seconds. Defaults to 86400 (24 hours); between 7200 and 2592000. Applies to future deliveries only.
Returns:

Successful response

def claim( self, subscription: str, input: EventSubscriptionClaimInput) -> archastro.platform.types.common.EventSubscriptionClaim:
347    def claim(
348        self, subscription: str, input: EventSubscriptionClaimInput
349    ) -> EventSubscriptionClaim:
350        """
351        Claim events from a subscription
352        Atomically leases the oldest unacknowledged delivery. Returns an empty data array when the queue is empty or its head already has an active lease. Passing max_events opts into batch mode, where leased entries are skipped instead of blocking; wait_seconds bounds a long poll on an empty queue.
353
354        Args:
355            subscription: Subscription ID (`esub_...`).
356            input: Request body.
357            input.consumer_id: Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.
358            input.max_events: Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.
359            input.request_id: Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.
360            input.visibility_timeout_seconds: Lease duration override between 15 and 3600 seconds.
361            input.wait_seconds: Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.
362
363        Returns:
364            Successful response
365        """
366        return self._http.request(
367            f"/api/v1/event_subscriptions/{subscription}/claim",
368            method="POST",
369            body=input,
370            response_type=EventSubscriptionClaim,
371        )

Claim events from a subscription Atomically leases the oldest unacknowledged delivery. Returns an empty data array when the queue is empty or its head already has an active lease. Passing max_events opts into batch mode, where leased entries are skipped instead of blocking; wait_seconds bounds a long poll on an empty queue.

Arguments:
  • subscription: Subscription ID (esub_...).
  • input: Request body.
  • input.consumer_id: Consumer identity between 1 and 128 bytes, recorded on the lease for attribution.
  • input.max_events: Opts into batch mode: leases up to this many lease-available deliveries (between 1 and 20) in sequence order, skipping leased entries instead of blocking on the head. Omit to keep strict head-of-line semantics.
  • input.request_id: Idempotency key between 1 and 128 bytes. Use a value unique per claim attempt (e.g. a UUID) the key is scoped to the subscription, so a reused value takes over whatever leases it last stamped. Retrying a claim with the same request_id while its leases are unexpired returns the same deliveries (regardless of max_events) with fresh receipt handles and refreshed leases.
  • input.visibility_timeout_seconds: Lease duration override between 15 and 3600 seconds.
  • input.wait_seconds: Long-poll bound between 0 and 20 seconds. When the queue is empty, the request waits up to this long for a delivery before returning an empty data array. 0 (or omitting) returns immediately.
Returns:

Successful response

def head( self, subscription: str) -> archastro.platform.types.common.EventSubscriptionHead:
373    def head(self, subscription: str) -> EventSubscriptionHead:
374        """
375        Peek at the head of a subscription queue
376        Returns the oldest unacknowledged delivery without reserving it. This diagnostic read cannot be used as a safe substitute for claim.
377
378        Args:
379            subscription: Subscription ID (`esub_...`).
380
381        Returns:
382            Successful response
383        """
384        return self._http.request(
385            f"/api/v1/event_subscriptions/{subscription}/head",
386            response_type=EventSubscriptionHead,
387        )

Peek at the head of a subscription queue Returns the oldest unacknowledged delivery without reserving it. This diagnostic read cannot be used as a safe substitute for claim.

Arguments:
  • subscription: Subscription ID (esub_...).
Returns:

Successful response

def queue( self, subscription: str, *, limit: int | None = None, before_cursor: str | None = None, after_cursor: str | None = None) -> archastro.platform.types.common.EventSubscriptionQueue:
389    def queue(
390        self,
391        subscription: str,
392        *,
393        limit: int | None = None,
394        before_cursor: str | None = None,
395        after_cursor: str | None = None,
396    ) -> EventSubscriptionQueue:
397        """
398        Read a subscription queue
399        Returns a non-reserving, oldest-first view of unacknowledged deliveries, including the queue's loss indicators (queue_epoch, dropped_events_total, dropped_through_position). Use claim to obtain a receipt handle before acknowledging.
400
401        Args:
402            subscription: Subscription ID (`esub_...`).
403            limit: Maximum entries to return, from 1 through 100.
404            before_cursor: Opaque cursor for the preceding page.
405            after_cursor: Opaque cursor for the following page.
406
407        Returns:
408            Successful response
409        """
410        query: dict[str, object] = {}
411        if limit is not None:
412            query["limit"] = limit
413        if before_cursor is not None:
414            query["before_cursor"] = before_cursor
415        if after_cursor is not None:
416            query["after_cursor"] = after_cursor
417        return self._http.request(
418            f"/api/v1/event_subscriptions/{subscription}/queue",
419            query=query,
420            response_type=EventSubscriptionQueue,
421        )

Read a subscription queue Returns a non-reserving, oldest-first view of unacknowledged deliveries, including the queue's loss indicators (queue_epoch, dropped_events_total, dropped_through_position). Use claim to obtain a receipt handle before acknowledging.

Arguments:
  • subscription: Subscription ID (esub_...).
  • limit: Maximum entries to return, from 1 through 100.
  • before_cursor: Opaque cursor for the preceding page.
  • after_cursor: Opaque cursor for the following page.
Returns:

Successful response