///|
let task_outcomes_applied_metadata_key : String = "posoco.task.applied_ids"
///|
let foreground_task_settle_timeout_ms : Int = 5_000
///|
let managed_task_shutdown_settle_timeout_ms : Int = 5_000
///|
/// Applied task ids may survive a session-store restart, so a process-local
/// counter alone cannot identify a new runtime. Use the platform entropy
/// source and fail activation if it is unavailable.
fn next_task_runtime_identity() -> String raise @error.AgentError {
let bytes = match @env.rand(16) {
Some(value) => value
None =>
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"task runtime nonce unavailable: platform entropy source returned no bytes",
),
)
}
let hex : Array[Char] = [
'0', '1', '2', '3', '4', '5', '6', '7', '8', '9', 'a', 'b', 'c', 'd', 'e', 'f',
]
let out = StringBuilder()
for i in 0..> 4) & 0x0f])
out.write_char(hex[byte & 0x0f])
}
out.to_string()
}
///|
priv struct AgentTaskOperation {
session_id : String
run_id : @kernel.RunId
turn_id : @kernel.TurnId
}
///|
priv struct AgentTaskRecord {
receipt : @port.TaskReceipt
task : @fuwaroid.SupervisedTask[@port.TaskOutcome]
operation : AgentTaskOperation?
}
///|
priv struct AgentTaskRuntime {
mut supervisor : @fuwaroid.Supervisor[@port.TaskOutcome]?
mut scope_active : Bool
mut closed : Bool
mut runtime_id : String?
mut next_id : Int
mut active_operation : AgentTaskOperation?
foreground_settle_timeout_ms : Int
shutdown_settle_timeout_ms : Int
tasks : Map[String, AgentTaskRecord]
outcomes : Map[String, Array[@port.TaskOutcome]]
inflight : Map[String, Array[String]]
}
///|
fn AgentTaskRuntime::AgentTaskRuntime(
foreground_settle_timeout_ms? : Int = foreground_task_settle_timeout_ms,
shutdown_settle_timeout_ms? : Int = managed_task_shutdown_settle_timeout_ms,
) -> AgentTaskRuntime {
{
supervisor: None,
scope_active: false,
closed: false,
runtime_id: None,
next_id: 1,
active_operation: None,
foreground_settle_timeout_ms,
shutdown_settle_timeout_ms,
tasks: Map::from_array([]),
outcomes: Map::from_array([]),
inflight: Map::from_array([]),
}
}
///|
fn AgentTaskRuntime::capability(
self : AgentTaskRuntime,
extension_id : String,
) -> @port.Tasks {
@port.Tasks::from_submit(submit=fn(spec) { self.submit(extension_id, spec) })
}
///|
fn[G] AgentTaskRuntime::activate(
self : AgentTaskRuntime,
group : @async.TaskGroup[G],
) -> Unit raise @error.AgentError {
if self.closed {
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed("task runtime is closed"),
)
}
if self.scope_active {
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed("task scope is already active"),
)
}
self.runtime_id = Some(next_task_runtime_identity())
self.supervisor = Some(@fuwaroid.Supervisor(group~))
self.scope_active = true
}
///|
fn AgentTaskRuntime::scope_available(self : AgentTaskRuntime) -> Bool {
!self.closed && !self.scope_active
}
///|
fn agent_task_operation_equal(
left : AgentTaskOperation,
right : AgentTaskOperation,
) -> Bool {
left.session_id == right.session_id &&
left.run_id == right.run_id &&
left.turn_id == right.turn_id
}
///|
fn AgentTaskRuntime::record_matches_operation(
_self : AgentTaskRuntime,
record : AgentTaskRecord,
operation : AgentTaskOperation,
) -> Bool {
match record.receipt.mode {
@port.TaskMode::Foreground =>
match record.operation {
Some(owner) => agent_task_operation_equal(owner, operation)
None => false
}
@port.TaskMode::Background => false
}
}
///|
fn supervised_task_terminal(status : @fuwaroid.SupervisedTaskStatus) -> Bool {
match status {
@fuwaroid.SupervisedTaskStatus::Completed
| @fuwaroid.SupervisedTaskStatus::Cancelled
| @fuwaroid.SupervisedTaskStatus::Failed => true
@fuwaroid.SupervisedTaskStatus::Running
| @fuwaroid.SupervisedTaskStatus::Cancelling => false
}
}
///|
fn AgentTaskRuntime::refresh_active_operation(self : AgentTaskRuntime) -> Unit {
match self.active_operation {
None => ()
Some(operation) => {
let stale : Array[String] = []
let mut unresolved = false
for record in self.tasks.values() {
if self.record_matches_operation(record, operation) {
if supervised_task_terminal(record.task.snapshot().status) {
stale.push(record.receipt.id)
} else {
unresolved = true
}
}
}
for id in stale {
self.reap(id)
}
if !unresolved {
self.active_operation = None
}
}
}
}
///|
fn AgentTaskRuntime::begin_operation(
self : AgentTaskRuntime,
session_id : String,
run_id : @kernel.RunId,
turn_id : @kernel.TurnId,
) -> Unit raise @error.AgentError {
self.refresh_active_operation()
if self.active_operation is Some(_) {
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"foreground task cleanup is still pending",
),
)
}
self.active_operation = Some({ session_id, run_id, turn_id, })
}
///|
fn AgentTaskRuntime::next_receipt(
self : AgentTaskRuntime,
extension_id : String,
spec : @port.TaskSpec,
) -> @port.TaskReceipt {
let runtime_id = match self.runtime_id {
Some(value) => value
None => abort("task runtime nonce missing while scope is active")
}
let id = "agent_task_\{runtime_id}_\{self.next_id}"
self.next_id = self.next_id + 1
{
id,
session_id: spec.session_id,
mode: spec.mode,
label: spec.label,
extension_id,
}
}
///|
fn AgentTaskRuntime::submit(
self : AgentTaskRuntime,
extension_id : String,
spec : @port.TaskSpec,
) -> Result[@port.TaskHandle, @port.TaskSubmitError] {
if spec.session_id == "" {
return Err(InvalidSession(reason="session_id must not be empty"))
}
match spec.timeout_ms {
Some(timeout_ms) if timeout_ms <= 0 =>
return Err(InvalidSession(reason="timeout_ms must be positive"))
_ => ()
}
if self.closed {
return Err(Closed)
}
let operation = self.active_operation
match spec.mode {
@port.TaskMode::Foreground =>
match operation {
None => return Err(Unavailable)
Some(owner) if owner.session_id != spec.session_id =>
return Err(
InvalidSession(
reason="foreground task session does not match the active operation",
),
)
Some(_) => ()
}
@port.TaskMode::Background => ()
}
let supervisor = match self.supervisor {
Some(value) => value
None => return Err(Unavailable)
}
let receipt = self.next_receipt(extension_id, spec)
let runtime = self
let task = match
supervisor.spawn(
async fn() { runtime.execute(spec, receipt, operation) },
label=receipt.label,
) {
Ok(value) => value
Err(_) => return Err(Closed)
}
self.tasks[receipt.id] = { receipt, task, operation, }
let handle_task = task
Ok(
@port.TaskHandle::from_callbacks(
receipt~,
wait=async fn(_target) {
let waited : Result[@port.TaskOutcome, Error] = Ok(handle_task.wait()) catch {
error => Err(error)
}
let outcome = match waited {
Ok(value) => value
Err(error) if error is @async.WaitedTaskAlreadyCancelled =>
@port.TaskOutcome::{
receipt,
status: @port.TaskStatus::Cancelled(reason="task cancelled"),
}
Err(error) =>
@port.TaskOutcome::{
receipt,
status: @port.TaskStatus::Failed(reason=error.to_string()),
}
}
runtime.reap(receipt.id)
snapshot_task_outcome(outcome)
},
cancel=fn(_target) { handle_task.cancel() },
),
)
}
///|
#warnings("-fragile_catch_all")
async fn AgentTaskRuntime::execute(
self : AgentTaskRuntime,
spec : @port.TaskSpec,
receipt : @port.TaskReceipt,
_operation : AgentTaskOperation?,
) -> @port.TaskOutcome {
let guarded : async () -> @kernel.Message = async fn() {
let message = (spec.run)() catch {
error => {
@async.pause()
raise error
}
}
@async.pause()
message
}
let status : @port.TaskStatus = match spec.timeout_ms {
Some(timeout_ms) => {
let result : Result[@kernel.Message?, Error] = Ok(
@async.handle_cancellation(async fn() {
@async.with_timeout(timeout_ms, () => guarded())
}),
) catch {
error => Err(error)
}
match result {
Ok(Some(message)) =>
@port.TaskStatus::Completed(message=snapshot_message(message))
Ok(None) => @port.TaskStatus::Cancelled(reason="task cancelled")
Err(error) if error is @async.TimeoutError => @port.TaskStatus::TimedOut
Err(error) => @port.TaskStatus::Failed(reason=error.to_string())
}
}
None => {
let result : Result[@kernel.Message?, Error] = Ok(
@async.handle_cancellation(guarded),
) catch {
error => Err(error)
}
match result {
Ok(Some(message)) =>
@port.TaskStatus::Completed(message=snapshot_message(message))
Ok(None) => @port.TaskStatus::Cancelled(reason="task cancelled")
Err(error) => @port.TaskStatus::Failed(reason=error.to_string())
}
}
}
let outcome : @port.TaskOutcome = { receipt, status, }
match receipt.mode {
@port.TaskMode::Background =>
self.enqueue_outcome(snapshot_task_outcome(outcome))
@port.TaskMode::Foreground => ()
}
outcome
}
///|
fn AgentTaskRuntime::reap(self : AgentTaskRuntime, receipt_id : String) -> Unit {
self.tasks.remove(receipt_id)
}
///|
fn AgentTaskRuntime::cancel_foreground(
self : AgentTaskRuntime,
operation : AgentTaskOperation,
) -> Unit {
for record in self.tasks.values() {
if self.record_matches_operation(record, operation) {
record.task.cancel()
}
}
}
///|
fn AgentTaskRuntime::cancel_foreground_for_run(
self : AgentTaskRuntime,
run_id : @kernel.RunId,
) -> Unit {
match self.active_operation {
Some(operation) if operation.run_id == run_id =>
self.cancel_foreground(operation)
_ => ()
}
}
///|
async fn AgentTaskRuntime::finish_active_operation(
self : AgentTaskRuntime,
) -> Unit raise @error.AgentError {
let operation = match self.active_operation {
None => return
Some(value) => value
}
let records = Array::from_iter(self.tasks.values()).filter(fn(record) {
self.record_matches_operation(record, operation)
})
if records.is_empty() {
self.active_operation = None
return
}
let supervisor = match self.supervisor {
Some(value) => value
None =>
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"task supervisor is unavailable during foreground cleanup",
),
)
}
let tasks = records.map(fn(record) { record.task })
let settled : Result[@fuwaroid.SettleOutcome, Error] = Ok(
@async.protect_from_cancel(async fn() {
supervisor.cancel_and_wait(
tasks=tasks.clamped_view(),
timeout_ms=self.foreground_settle_timeout_ms,
)
}),
) catch {
error => Err(error)
}
let outcome = match settled {
Ok(value) => value
Err(error) =>
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"foreground task settlement failed: " + error.to_string(),
),
)
}
match outcome {
@fuwaroid.Settled => {
for record in records {
self.reap(record.receipt.id)
}
self.active_operation = None
}
@fuwaroid.DeadlineExceeded(snapshots) => {
self.refresh_active_operation()
let detail = snapshots
.map(fn(snapshot) {
let label = if snapshot.label == "" {
""
} else {
snapshot.label
}
"\{label}:\{repr(snapshot.status)}"
})
.join(", ")
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"foreground task cleanup deadline exceeded: " + detail,
),
)
}
}
}
///|
fn AgentTaskRuntime::enqueue_outcome(
self : AgentTaskRuntime,
outcome : @port.TaskOutcome,
) -> Unit {
let session_id = outcome.receipt.session_id
if self.outcomes.contains(session_id) {
let existing = self.outcomes[session_id]
if !existing.iter().any(fn(item) { item.receipt.id == outcome.receipt.id }) {
existing.push(outcome)
}
} else {
self.outcomes[session_id] = [outcome]
}
}
///|
fn AgentTaskRuntime::has_inflight(
self : AgentTaskRuntime,
session_id : String,
) -> Bool {
match self.inflight.get(session_id) {
Some(ids) => !ids.is_empty()
None => false
}
}
///|
fn AgentTaskRuntime::metadata_with_applied_outcomes(
self : AgentTaskRuntime,
session_id : String,
metadata : Map[String, Json],
) -> Map[String, Json] raise @error.AgentError {
let ids = parse_task_applied_ids(metadata)
match self.inflight.get(session_id) {
Some(inflight) =>
for id in inflight {
if !ids.contains(id) {
ids.push(id)
}
}
None => ()
}
if !ids.is_empty() {
metadata[task_outcomes_applied_metadata_key] = Json::array(
ids.map(fn(id) { Json::string(id) }),
)
}
metadata
}
///|
/// Parse the Agent-owned applied-outcome metadata once at each session
/// admission. Any malformed value is a runtime error before the turn can
/// perform model, memory, or task-outcome side effects.
fn parse_task_applied_ids(
metadata : Map[String, Json],
) -> Array[String] raise @error.AgentError {
match metadata.get(task_outcomes_applied_metadata_key) {
None => []
Some(Array(values)) => {
let ids : Array[String] = []
for value in values {
match value {
String(id) => ids.push(id)
_ =>
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"malformed posoco.task.applied_ids metadata entry",
),
)
}
}
ids
}
Some(_) =>
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"malformed posoco.task.applied_ids metadata value",
),
)
}
}
///|
fn task_outcome_is_applied(
applied_ids : Array[String],
outcome : @port.TaskOutcome,
) -> Bool {
applied_ids.contains(outcome.receipt.id)
}
///|
fn AgentTaskRuntime::peek_outcomes(
self : AgentTaskRuntime,
session_id : String,
) -> Array[@port.TaskOutcome] {
let pending : Array[@port.TaskOutcome] = match self.outcomes.get(session_id) {
Some(values) => values.map(snapshot_task_outcome)
None => []
}
self.inflight[session_id] = pending.map(fn(outcome) { outcome.receipt.id })
pending
}
///|
/// Copy the message payload carried by a task terminal so queued delivery and
/// repeated handle waits cannot share mutable content with the worker result.
fn snapshot_task_outcome(outcome : @port.TaskOutcome) -> @port.TaskOutcome {
let status = match outcome.status {
@port.TaskStatus::Completed(message~) =>
@port.TaskStatus::Completed(message=snapshot_message(message))
@port.TaskStatus::Failed(reason~) => @port.TaskStatus::Failed(reason~)
@port.TaskStatus::TimedOut => @port.TaskStatus::TimedOut
@port.TaskStatus::Cancelled(reason~) => @port.TaskStatus::Cancelled(reason~)
}
{ receipt: outcome.receipt, status, }
}
///|
fn AgentTaskRuntime::commit_outcomes(
self : AgentTaskRuntime,
session_id : String,
) -> Unit {
let committed : Array[String] = match self.inflight.get(session_id) {
Some(ids) => ids
None => return
}
match self.outcomes.get(session_id) {
Some(values) => {
let remaining = values.filter(fn(outcome) {
!committed.contains(outcome.receipt.id)
})
if remaining.is_empty() {
self.outcomes.remove(session_id)
} else {
self.outcomes[session_id] = remaining
}
for id in committed {
self.reap(id)
}
}
None => ()
}
self.inflight.remove(session_id)
}
///|
fn task_message_text(message : @kernel.Message) -> String {
let content_text = fn(content : @kernel.Content) -> String {
match content {
@kernel.Content::Text(text) => text
@kernel.Content::Image(media_type~, ..) => "[image=\{media_type}]"
}
}
match message {
@kernel.Message::SystemMessage(content~) =>
content.map(content_text).join("")
@kernel.Message::UserMessage(content~) => content.map(content_text).join("")
@kernel.Message::AssistantMessage(content~, ..) =>
content.map(content_text).join("")
@kernel.Message::ToolMessage(outcome~, ..) => outcome.summary()
}
}
///|
fn task_outcome_marker(outcome : @port.TaskOutcome) -> String {
"[posoco-task-outcome id=\{outcome.receipt.id}]"
}
///|
/// Outcomes enter the transcript as user/context data. This prevents a task
/// completion from gaining system or tool authority merely because it came
/// from an Agent-owned worker.
fn task_outcome_message(outcome : @port.TaskOutcome) -> @kernel.Message {
let marker = task_outcome_marker(outcome)
let detail = match outcome.status {
@port.TaskStatus::Completed(message~) => task_message_text(message)
@port.TaskStatus::Failed(reason~) => "failed: \{reason}"
@port.TaskStatus::TimedOut => "timed out"
@port.TaskStatus::Cancelled(reason~) => "cancelled: \{reason}"
}
@kernel.Message::UserMessage(content=[
@kernel.Content::Text(
"\{marker} extension=\{outcome.receipt.extension_id} label=\{outcome.receipt.label}: \{detail}",
),
])
}
///|
fn AgentTaskRuntime::shutdown_straggler_detail(
self : AgentTaskRuntime,
) -> String {
let details : Array[String] = []
for record in self.tasks.values() {
let snapshot = record.task.snapshot()
if !supervised_task_terminal(snapshot.status) {
let label = if record.receipt.label == "" {
""
} else {
record.receipt.label
}
details.push(
"extension=\{record.receipt.extension_id} label=\{label} receipt=\{record.receipt.id} status=\{repr(snapshot.status)}",
)
}
}
if details.is_empty() {
""
} else {
details.join(", ")
}
}
///|
/// Close task admission and bounded-settle all live managed work.
///
/// The deadline bounds Posoco's settlement observation only. Managed tasks
/// remain children of the owning async TaskGroup, so a worker that ignores
/// cooperative cancellation can still keep the structured scope alive after
/// this method reports the offending task.
async fn AgentTaskRuntime::shutdown(
self : AgentTaskRuntime,
) -> Unit raise @error.AgentError {
if !self.closed {
self.closed = true
self.scope_active = false
}
let supervisor = match self.supervisor {
Some(value) => value
None => {
self.active_operation = None
return
}
}
let settled : Result[@fuwaroid.SettleOutcome, Error] = Ok(
@async.protect_from_cancel(async fn() {
supervisor.shutdown(timeout_ms=self.shutdown_settle_timeout_ms)
}),
) catch {
error => Err(error)
}
let outcome = match settled {
Ok(value) => value
Err(error) =>
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"managed task shutdown settlement failed: " + error.to_string(),
),
)
}
match outcome {
@fuwaroid.Settled => {
let records : Array[AgentTaskRecord] = Array::from_iter(
self.tasks.values(),
)
for record in records {
self.reap(record.receipt.id)
}
self.supervisor = None
self.active_operation = None
}
@fuwaroid.DeadlineExceeded(_) =>
raise @error.AgentError::Runtime(
@error.RuntimeError::InvocationFailed(
"managed task shutdown deadline exceeded: " +
self.shutdown_straggler_detail(),
),
)
}
}