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
|
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 |
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.
as_actor
¶
Return a handle on the same workspace that acts as actor.
commit
async
¶
commit(
command: CreateArtifact
| EditArtifact
| ArchiveArtifact,
) -> Applied | Proposed
commit(command: ProposeChange) -> Proposed
commit(command: RespondToProposal) -> Resolved
commit(
command: CreateThread
| PostMessage
| SetFocus
| SetThreadMode
| AnswerDeferred
| GiveFeedback,
) -> Recorded
record
async
¶
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.
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 |
None
|
kind
|
MessageKind
|
|
'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
¶
proposal
async
¶
proposal(proposal_id: ProposalId) -> 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
¶
runs
async
¶
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.
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 |
0
|
before_seq
|
int | None
|
Only envelopes before this |
None
|
threads
|
Collection[ThreadId] | None
|
Only thread-scoped events from these threads; workspace-scoped events
(artifacts, proposals) are always included. |
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 |
None
|
Raises:
| Type | Description |
|---|---|
ValidationFailed
|
If both |
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 |
required |
before_seq
|
int | None
|
Only consider envelopes with a lesser |
None
|
focus
|
Collection[ArtifactId] | None
|
The artifacts the viewer follows; |
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
¶
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 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 |
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. |
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.
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.
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.
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 |
0
|
before_seq
|
int | None
|
Only envelopes before this |
None
|
threads
|
Collection[ThreadId] | None
|
Only the envelopes a subscriber following these threads receives, by
|
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
¶
Take or renew an exclusive lease. Return False if another holder has it.
release_lease
async
¶
Release a lease if holder has it.
cursor
async
¶
Return how far a named consumer of the log has got: the seq saved, or 0.
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 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 |
Scope
dataclass
¶
Scope(tenant_id: TenantId, workspace_id: WorkspaceId)
A tenant's workspace: the unit of isolation, ordering and locking.
HistoryChunk
dataclass
¶
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
¶
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.
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.
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.
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
¶
Take or renew an exclusive lease. Return False if another holder has it.
release_lease
async
¶
Release a lease if holder has it.
cursor
async
¶
Return how far a named consumer of the log has got: the seq saved, or 0.