Skip to content

artifactr.workspace

Tenant-scoped workspaces over pluggable storage.

Open a Workspace with Workspaces.open; every read and write goes through it, and every write runs core's rules in one storage transaction.

Workspaces

Tenant-scoped handles. See Workspaces, commits and the log.

Workspaces

Workspaces(
    storage: Storage,
    *,
    types: Iterable[type[Artifact]] | None = None,
    tracer_provider: TracerProvider | None = None,
    meter_provider: MeterProvider | None = None,
)

Opens scoped Workspace handles over one storage.

Parameters:

Name Type Description Default
storage Storage

Where workspaces are kept.

required
types Iterable[type[Artifact]] | None

The artifact types this application accepts. Others are rejected even if they are registered, so clients cannot create arbitrary types. None accepts every registered type.

None
tracer_provider TracerProvider | None

Where commit spans go. Defaults to the global tracer provider.

None
meter_provider MeterProvider | None

Where artifactr's metrics go. Defaults to the global meter provider.

None

open async

open(
    tenant_id: TenantId,
    workspace_id: WorkspaceId,
    *,
    actor: Actor,
    authorize: Authorize | None = None,
) -> Workspace

Return a handle on a tenant's workspace, acting as actor.

This is the only place a tenant id enters; nothing on the handle can reach another tenant. The surfaces pass their authorize hook, so each refuses a workspace alike.

Raises:

Type Description
Forbidden

If authorize is given and refuses the actor this workspace.

Workspace

Workspace(
    storage: Storage,
    scope: Scope,
    actor: Actor,
    kinds: frozenset[str] | None,
    telemetry: Telemetry | None = None,
)

A handle on one tenant's workspace, bound to the actor it acts as.

Create handles with Workspaces.open, and derive handles for other actors (such as the agent) with as_actor.

tenant_id property

tenant_id: TenantId

The id of the tenant the workspace belongs to.

workspace_id property

workspace_id: WorkspaceId

The workspace's id.

actor property

actor: Actor

Who this handle acts as.

as_actor

as_actor(actor: Actor) -> Workspace

Return a handle on the same workspace that acts as actor.

commit async

commit(command: Command) -> Outcome

Submit a command.

Returns:

Type Description
Outcome

What the command did, with the seq of the last event it appended.

Raises:

Type Description
Rejection

If the command cannot be applied.

record async

record(
    fact: Fact, *, history: bytes | None = None
) -> Recorded

Record a fact about an agent run, or an application event.

Parameters:

Name Type Description Default
fact Fact

The fact to record.

required
history bytes | None

Serialized model messages to append to the fact's thread history in the same transaction, typically with RunPaused or RunEnded.

None

create async

create(
    artifact: Artifact,
    *,
    artifact_id: ArtifactId | None = None,
    thread_id: ThreadId | None = None,
) -> Applied | Proposed

Create an artifact from an instance of its type.

create_thread async

create_thread(title: str = '') -> Thread

Create a thread and return it.

post_message async

post_message(
    thread_id: ThreadId,
    content: str,
    *,
    message_id: MessageId | None = None,
    kind: MessageKind = "message",
) -> Recorded

Post a message, or a notice, in a thread as this handle's actor.

Parameters:

Name Type Description Default
thread_id ThreadId

The thread to post in.

required
content str

What the message says.

required
message_id MessageId | None

The message's id; a new one when omitted. An id already used in the workspace is refused with InvalidState.

None
kind MessageKind

notice for a notice, which is for people and starts no run (ADR-0051).

'message'

get async

get(
    artifact_type: type[A], artifact_id: ArtifactId
) -> Versioned[A]

Return an artifact's current version, checked to be of artifact_type.

Raises:

Type Description
NotFound

If there is no such artifact of that type.

artifact async

artifact(artifact_id: ArtifactId) -> Versioned[Artifact]

Return an artifact's current version, whatever its type.

Raises:

Type Description
NotFound

If there is no such artifact.

artifacts async

artifacts(
    artifact_type: type[A] = Artifact,
    *,
    kind: str | None = None,
    include_archived: bool = False,
) -> list[Versioned[A]]

Return current artifacts that are instances of artifact_type, oldest first.

Parameters:

Name Type Description Default
artifact_type type[A]

Only artifacts of this type, or a subclass of it.

Artifact
kind str | None

Only artifacts registered under this kind.

None
include_archived bool

Archived artifacts too.

False

revisions async

revisions(artifact_id: ArtifactId) -> list[Revision]

Return an artifact's revisions, oldest first.

Raises:

Type Description
NotFound

If there is no such artifact: every artifact has a revision.

thread async

thread(thread_id: ThreadId) -> Thread

Return a thread.

Raises:

Type Description
NotFound

If there is no such thread.

threads async

threads() -> list[Thread]

Return every thread, oldest first.

proposal async

proposal(proposal_id: ProposalId) -> Proposal

Return a proposal.

Raises:

Type Description
NotFound

If there is no such proposal.

proposals async

proposals(
    *,
    status: Literal["pending", "accepted", "rejected"]
    | None = "pending",
) -> list[Proposal]

Return proposals with a status (pending by default; None for all), oldest first.

run async

run(run_id: RunId) -> Run

Return a run.

Raises:

Type Description
NotFound

If there is no such run.

runs async

runs(
    *,
    thread_id: ThreadId | None = None,
    status: RunStatus | None = None,
) -> list[Run]

Return runs, optionally of one thread and with one status, oldest first.

history async

history(thread_id: ThreadId) -> Sequence[HistoryChunk]

Return a thread's history: serialized model messages, one chunk per run segment.

head_seq async

head_seq() -> int

Return the log's latest seq, or 0 if it is empty.

read async

read(
    *,
    after_seq: int = 0,
    before_seq: int | None = None,
    threads: Collection[ThreadId] | None = None,
    limit: int | None = None,
    last: int | None = None,
) -> list[Envelope]

Return logged envelopes in the window after_seq < seq < before_seq, in order.

Parameters:

Name Type Description Default
after_seq int

Only envelopes after this seq.

0
before_seq int | None

Only envelopes before this seq; None reads to the head.

None
threads Collection[ThreadId] | None

Only thread-scoped events from these threads; workspace-scoped events (artifacts, proposals) are always included. None includes every thread.

None
limit int | None

At most this many: the first ones in the window that match.

None
last int | None

At most this many: the last ones in the window that match, still in order. To page backwards, read last=n, then last=n before the oldest seq returned.

None

Raises:

Type Description
ValidationFailed

If both limit and last are given, or a number is negative.

subscribe async

subscribe(
    *,
    after_seq: int = 0,
    threads: Collection[ThreadId] | None = None,
) -> AsyncIterator[Envelope]

Yield envelopes after after_seq: the stored ones, then new ones as they commit.

Replay and live delivery are the same stream, so nothing falls between them. Reading the log is untraced, so a subscription that polls storage makes no trace per poll (ADR-0046); what the subscriber does with each envelope is traced as usual.

change_notes async

change_notes(
    *,
    after_seq: int,
    before_seq: int | None = None,
    focus: Collection[ArtifactId] | None = None,
    viewer: Actor | None = None,
    notices: bool = False,
) -> list[Note]

Return notes about what others did in the window after_seq < seq < before_seq.

Parameters:

Name Type Description Default
after_seq int

Only consider envelopes with a greater seq.

required
before_seq int | None

Only consider envelopes with a lesser seq; None reads to the head.

None
focus Collection[ArtifactId] | None

The artifacts the viewer follows; None means all.

None
viewer Actor | None

Who the notes are for; defaults to this handle's actor.

None
notices bool

Include the notices others posted in the viewer's thread, when the viewer is a thread's agent.

False

cursor async

cursor(name: str) -> int

Return how far a named consumer of the log has got: the seq it saved, or 0.

A consumer, such as a feedback mirror, saves its cursor with save_cursor and carries on after it when it starts again.

save_cursor async

save_cursor(name: str, seq: int) -> None

Save how far a named consumer of the log has got: it is done with seq.

A cursor only moves forward: saving a seq below the saved one leaves it, so a consumer that runs in several processes, each at its own pace, cannot move it back.

claim_thread async

claim_thread(
    thread_id: ThreadId,
    *,
    holder: str,
    ttl: timedelta = timedelta(seconds=30),
) -> AsyncGenerator[Event]

Hold a thread exclusively, renewing the claim until the block exits.

One run is active per thread. The claim is a lease in storage, so it holds across processes and lapses by itself if the holder dies.

Yields:

Type Description
AsyncGenerator[Event]

An event set once the claim is lost: it lapsed, a whole ttl after its last

AsyncGenerator[Event]

renewal began, even while a renewal hangs, or a renewal found that another holder

AsyncGenerator[Event]

has the thread. The claim is never renewed after that, even if storage would allow

AsyncGenerator[Event]

it. Its lease is released only when the block exits, so exit promptly once the event

AsyncGenerator[Event]

is set.

Raises:

Type Description
ThreadBusy

If another holder has the thread.

ThreadBusy

ThreadBusy(message: str)

Bases: InvalidState

Another run holds the thread.

Authorize module-attribute

Decides whether an actor may use a workspace of its tenant.

The surfaces that serve workspaces, the FastAPI router and the MCP server, take one as their authorize hook, and pass it to Workspaces.open for every request that names a workspace.

Storage

The storage protocol, and the in-memory implementation. See Storage.

Storage

Bases: Protocol

Persistence for workspaces. Every method is scoped to one tenant's workspace.

A caller can be cancelled at any await, by a stopped run, a timeout or a closed connection, so every method is cancel-safe: a cancelled caller leaves no lock or connection behind; cancellation may be deferred until the current statement ends. So a cancelled call may still have taken effect: a transaction may have committed, or acquire_lease taken the lease. A wait for a pooled connection cannot be interrupted either, so a task that holds a transaction must not await a task it cancelled. The event loop's shutdown is not covered: it cancels the storage's own tasks too, so stop runs with Runner.aclose first.

transaction

transaction(
    scope: Scope,
) -> AbstractAsyncContextManager[Transaction]

Begin a transaction that commits when the block exits normally.

artifact async

artifact(
    scope: Scope, artifact_id: ArtifactId
) -> Versioned[Artifact] | None

Return an artifact's current version, or None.

artifacts async

artifacts(
    scope: Scope,
    *,
    kind: str | None = None,
    include_archived: bool = False,
) -> list[Versioned[Artifact]]

Return current artifacts, optionally of one kind, oldest first.

revisions async

revisions(
    scope: Scope, artifact_id: ArtifactId
) -> list[Revision]

Return an artifact's revisions, oldest first.

thread async

thread(scope: Scope, thread_id: ThreadId) -> Thread | None

Return a thread, or None.

threads async

threads(scope: Scope) -> list[Thread]

Return every thread, oldest first.

proposal async

proposal(
    scope: Scope, proposal_id: ProposalId
) -> Proposal | None

Return a proposal, or None.

proposals async

proposals(
    scope: Scope,
    *,
    status: Literal["pending", "accepted", "rejected"]
    | None = None,
) -> list[Proposal]

Return proposals, optionally with one status, oldest first.

run async

run(scope: Scope, run_id: RunId) -> Run | None

Return a run, or None.

runs async

runs(
    scope: Scope,
    *,
    thread_id: ThreadId | None = None,
    status: RunStatus | None = None,
) -> list[Run]

Return runs, optionally of one thread and with one status, oldest first.

head_seq async

head_seq(scope: Scope) -> int

Return the log's latest seq, or 0 if it is empty.

read async

read(
    scope: Scope,
    *,
    after_seq: int = 0,
    before_seq: int | None = None,
    threads: Collection[ThreadId] | None = None,
    limit: int | None = None,
    last: int | None = None,
) -> list[Envelope]

Return logged envelopes in the window after_seq < seq < before_seq, in order.

Parameters:

Name Type Description Default
scope Scope

The workspace whose log to read.

required
after_seq int

Only envelopes after this seq.

0
before_seq int | None

Only envelopes before this seq; None reads to the head.

None
threads Collection[ThreadId] | None

Only the envelopes a subscriber following these threads receives, by delivered_to's rule; None reads every thread.

None
limit int | None

At most this many: the first ones in the window that match.

None
last int | None

At most this many: the last ones in the window that match, still returned in order. This is the tail of the log.

None

Callers give at most one of limit and last, and no negative numbers; read checks.

subscribe

subscribe(
    scope: Scope, *, after_seq: int = 0
) -> AsyncGenerator[Envelope]

Yield envelopes with seq greater than after_seq.

First the stored ones, then each new one as it is committed. The iterator runs until it is closed.

history async

history(
    scope: Scope, thread_id: ThreadId
) -> Sequence[HistoryChunk]

Return a thread's history chunks, in the order they were appended.

acquire_lease async

acquire_lease(
    scope: Scope, key: str, holder: str, ttl: timedelta
) -> bool

Take or renew an exclusive lease. Return False if another holder has it.

release_lease async

release_lease(scope: Scope, key: str, holder: str) -> None

Release a lease if holder has it.

cursor async

cursor(scope: Scope, name: str) -> int

Return how far a named consumer of the log has got: the seq saved, or 0.

save_cursor async

save_cursor(scope: Scope, name: str, seq: int) -> None

Save how far a named consumer of the log has got, such as a feedback mirror.

A cursor only moves forward: saving a seq below the saved one leaves it, so a consumer running in several processes cannot move it back.

Transaction

Bases: Protocol

One atomic unit of work on a workspace.

A transaction must be serialized with respect to every other transaction on the same scope from the moment it begins, so what it loads cannot change before it saves. The in-memory storage holds a per-workspace lock; SQL storage locks the workspace row.

load async

load(needs: Needs) -> State

Load the requested entities, and whether the requested message ids are used.

Entity ids that do not exist map to None.

save async

save(
    result: CommitResult,
    *,
    actor: Actor,
    traceparent: str | None = None,
) -> list[Envelope]

Persist a result's entities, revisions and used message ids, and log its events.

Parameters:

Name Type Description Default
result CommitResult

What core decided.

required
actor Actor

Who the envelopes are attributed to.

required
traceparent str | None

The W3C trace context the envelopes record.

None

Returns:

Type Description
list[Envelope]

The appended envelopes, with their seq numbers assigned.

append_history async

append_history(
    thread_id: ThreadId, messages: bytes
) -> None

Append serialized model messages to a thread's history.

Scope dataclass

Scope(tenant_id: TenantId, workspace_id: WorkspaceId)

A tenant's workspace: the unit of isolation, ordering and locking.

HistoryChunk dataclass

HistoryChunk(seq: int, messages: bytes)

Serialized model messages from one run segment, and the seq they were saved at.

seq instance-attribute

seq: int

The log's head when the chunk was saved: everything up to it happened before.

seal

seal(
    events: Sequence[KnownEvent],
    *,
    after_seq: int,
    scope: Scope,
    actor: Actor,
    ts: datetime,
    traceparent: str | None,
) -> list[Envelope]

Put a result's events in envelopes, numbered on from after_seq, to log them.

Implementations of Transaction.save call it, so every storage stamps, scopes and numbers envelopes alike.

InMemoryStorage

InMemoryStorage(*, clock: Clock = _utc_now)

Storage that keeps every workspace in process memory.

Parameters:

Name Type Description Default
clock Clock

Returns the current time. Defaults to the system clock in UTC.

_utc_now

transaction async

transaction(scope: Scope) -> AsyncGenerator[_Transaction]

Begin a transaction; it holds the workspace's lock until it ends.

artifact async

artifact(
    scope: Scope, artifact_id: ArtifactId
) -> Versioned[Artifact] | None

Return an artifact's current version, or None.

artifacts async

artifacts(
    scope: Scope,
    *,
    kind: str | None = None,
    include_archived: bool = False,
) -> list[Versioned[Artifact]]

Return current artifacts, optionally of one kind, oldest first.

revisions async

revisions(
    scope: Scope, artifact_id: ArtifactId
) -> list[Revision]

Return an artifact's revisions, oldest first.

thread async

thread(scope: Scope, thread_id: ThreadId) -> Thread | None

Return a thread, or None.

threads async

threads(scope: Scope) -> list[Thread]

Return every thread, oldest first.

proposal async

proposal(
    scope: Scope, proposal_id: ProposalId
) -> Proposal | None

Return a proposal, or None.

proposals async

proposals(
    scope: Scope,
    *,
    status: Literal["pending", "accepted", "rejected"]
    | None = None,
) -> list[Proposal]

Return proposals, optionally with one status, oldest first.

run async

run(scope: Scope, run_id: RunId) -> Run | None

Return a run, or None.

runs async

runs(
    scope: Scope,
    *,
    thread_id: ThreadId | None = None,
    status: RunStatus | None = None,
) -> list[Run]

Return runs, optionally of one thread and with one status, oldest first.

head_seq async

head_seq(scope: Scope) -> int

Return the log's latest seq, or 0 if it is empty.

read async

read(
    scope: Scope,
    *,
    after_seq: int = 0,
    before_seq: int | None = None,
    threads: Collection[ThreadId] | None = None,
    limit: int | None = None,
    last: int | None = None,
) -> list[Envelope]

Return logged envelopes in the window after_seq < seq < before_seq, in order.

Of those delivered to threads, if given: the first limit or the last last.

subscribe async

subscribe(
    scope: Scope, *, after_seq: int = 0
) -> AsyncGenerator[Envelope]

Yield stored envelopes after after_seq, then each new one as it commits.

history async

history(
    scope: Scope, thread_id: ThreadId
) -> Sequence[HistoryChunk]

Return a thread's history chunks, in the order they were appended.

acquire_lease async

acquire_lease(
    scope: Scope, key: str, holder: str, ttl: timedelta
) -> bool

Take or renew an exclusive lease. Return False if another holder has it.

release_lease async

release_lease(scope: Scope, key: str, holder: str) -> None

Release a lease if holder has it.

cursor async

cursor(scope: Scope, name: str) -> int

Return how far a named consumer of the log has got: the seq saved, or 0.

save_cursor async

save_cursor(scope: Scope, name: str, seq: int) -> None

Save how far a named consumer of the log has got; a cursor only moves forward.