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,
)

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

open async

open(
    tenant_id: TenantId,
    workspace_id: WorkspaceId,
    *,
    actor: Actor,
) -> 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.

Workspace

Workspace(
    storage: Storage,
    scope: Scope,
    actor: Actor,
    kinds: frozenset[str] | 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.

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: ProposeChange) -> Proposed
commit(command: RespondToProposal) -> Resolved
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) -> Recorded

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

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,
    *,
    include_archived: bool = False,
) -> list[Versioned[A]]

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

revisions async

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

Return an artifact's revisions, oldest first.

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,
    threads: Collection[ThreadId] | None = None,
    limit: int | None = None,
) -> list[Envelope]

Return logged envelopes after after_seq, in order.

Parameters:

Name Type Description Default
after_seq int

Return envelopes with a greater seq.

0
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

The most envelopes to return.

None

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.

change_notes async

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

Return notes about what others did since after_seq.

Parameters:

Name Type Description Default
after_seq int

Only consider envelopes with a greater seq.

required
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

claim_thread async

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

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.

Raises:

Type Description
ThreadBusy

If another holder has the thread.

ThreadBusy

ThreadBusy(message: str)

Bases: InvalidState

Another run holds the thread.

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.

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,
    limit: int | None = None,
) -> list[Envelope]

Return logged envelopes with seq greater than after_seq, in order.

subscribe

subscribe(
    scope: Scope, *, after_seq: int = 0
) -> AsyncIterator[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.

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. Ids that do not exist map to None.

save async

save(
    result: CommitResult, *, actor: Actor
) -> list[Envelope]

Persist a result's entities and revisions, and append its events to the log.

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.

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,
    limit: int | None = None,
) -> list[Envelope]

Return logged envelopes with seq greater than after_seq, in order.

subscribe async

subscribe(
    scope: Scope, *, after_seq: int = 0
) -> AsyncIterator[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.