Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 10 additions & 8 deletions harness/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions harness/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,10 @@ name = "harness"
path = "src/lib.rs"

[dependencies]
iii-sdk = "=0.19.2"
iii-sdk = "=0.20.0"
# Re-exports the same `opentelemetry` the SDK uses, so we can carry the active
# OTel context across `tokio::spawn` boundaries (shared `Context` global).
iii-observability = "0.19.2"
iii-helpers = "=0.20.0"
tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "signal", "time"] }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
Expand Down
7 changes: 4 additions & 3 deletions harness/src/clients/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,8 @@

use std::sync::Arc;

use iii_sdk::{TriggerRequest, III};
use iii_sdk::protocol::TriggerRequest;
use iii_sdk::IIIClient;
use serde::Deserialize;
use serde_json::{json, Value};

Expand Down Expand Up @@ -44,12 +45,12 @@ pub struct AssembleParams {

#[derive(Clone)]
pub struct ContextClient {
iii: Arc<III>,
iii: Arc<IIIClient>,
timeout_ms: u64,
}

impl ContextClient {
pub fn new(iii: Arc<III>, timeout_ms: u64) -> Self {
pub fn new(iii: Arc<IIIClient>, timeout_ms: u64) -> Self {
Self { iii, timeout_ms }
}

Expand Down
7 changes: 4 additions & 3 deletions harness/src/clients/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@

use std::sync::Arc;

use iii_sdk::{TriggerRequest, III};
use iii_sdk::protocol::TriggerRequest;
use iii_sdk::IIIClient;
use serde_json::{json, Value};

use crate::error::HarnessError;
Expand All @@ -21,12 +22,12 @@ pub struct FunctionDescriptor {

#[derive(Clone)]
pub struct EngineClient {
iii: Arc<III>,
iii: Arc<IIIClient>,
timeout_ms: u64,
}

impl EngineClient {
pub fn new(iii: Arc<III>, timeout_ms: u64) -> Self {
pub fn new(iii: Arc<IIIClient>, timeout_ms: u64) -> Self {
Self { iii, timeout_ms }
}

Expand Down
11 changes: 6 additions & 5 deletions harness/src/clients/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,10 @@ use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};

use async_trait::async_trait;
use iii_observability::opentelemetry::trace::FutureExt as _;
use iii_helpers::observability::opentelemetry::trace::FutureExt as _;
use iii_sdk::helpers::create_channel;
use iii_sdk::{TriggerRequest, III};
use iii_sdk::protocol::TriggerRequest;
use iii_sdk::IIIClient;
use serde_json::{json, Value};
use tokio::sync::mpsc;
use tokio::sync::Notify;
Expand Down Expand Up @@ -53,13 +54,13 @@ pub struct ChatOutcome {

#[derive(Clone)]
pub struct RouterClient {
iii: Arc<III>,
iii: Arc<IIIClient>,
timeout_ms: u64,
coalesce_ms: u64,
}

impl RouterClient {
pub fn new(iii: Arc<III>, timeout_ms: u64, coalesce_ms: u64) -> Self {
pub fn new(iii: Arc<IIIClient>, timeout_ms: u64, coalesce_ms: u64) -> Self {
Self {
iii,
timeout_ms,
Expand Down Expand Up @@ -141,7 +142,7 @@ impl RouterClient {
// the SDK uses to attach a handler's context.
let iii = self.iii.clone();
let timeout_ms = self.timeout_ms;
let parent_cx = iii_observability::opentelemetry::Context::current();
let parent_cx = iii_helpers::observability::opentelemetry::Context::current();
let trigger = tokio::spawn(
async move {
iii.trigger(TriggerRequest {
Expand Down
7 changes: 4 additions & 3 deletions harness/src/clients/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,8 @@

use std::sync::Arc;

use iii_sdk::{TriggerRequest, III};
use iii_sdk::protocol::TriggerRequest;
use iii_sdk::IIIClient;
use serde::Deserialize;
use serde_json::{json, Value};

Expand All @@ -30,12 +31,12 @@ pub struct LoadedCustom {

#[derive(Clone)]
pub struct SessionClient {
iii: Arc<III>,
iii: Arc<IIIClient>,
timeout_ms: u64,
}

impl SessionClient {
pub fn new(iii: Arc<III>, timeout_ms: u64) -> Self {
pub fn new(iii: Arc<IIIClient>, timeout_ms: u64) -> Self {
Self { iii, timeout_ms }
}

Expand Down
35 changes: 21 additions & 14 deletions harness/src/configuration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,10 @@
use std::sync::Arc;
use std::time::Duration;

use iii_sdk::{IIIError, RegisterFunction, RegisterTriggerInput, Trigger, TriggerRequest, III};
use iii_sdk::errors::Error;
use iii_sdk::protocol::{RegisterTriggerInput, TriggerRequest};
use iii_sdk::trigger::Trigger;
use iii_sdk::{IIIClient, RegisterFunction};
use serde_json::{json, Value};
use tokio::sync::RwLock;

Expand All @@ -38,7 +41,7 @@ const CONFIG_RETRY_BACKOFF_MS: u64 = 250;
/// Register the `harness` configuration schema. When `seed` is present its
/// value is installed as `initial_value`; otherwise the built-in default is
/// seeded only when nothing is stored yet (safe to call every boot).
pub async fn register_config(iii: &III, seed: Option<&WorkerConfig>) -> Result<(), String> {
pub async fn register_config(iii: &IIIClient, seed: Option<&WorkerConfig>) -> Result<(), String> {
discard_legacy_harness_config_if_needed(iii).await?;

let mut payload = json!({
Expand All @@ -61,7 +64,7 @@ pub async fn register_config(iii: &III, seed: Option<&WorkerConfig>) -> Result<(

/// Read the live `harness` configuration (env-expanded by the configuration
/// worker — `from_json` does NOT re-expand).
pub async fn fetch_config(iii: &III) -> Result<WorkerConfig, String> {
pub async fn fetch_config(iii: &IIIClient) -> Result<WorkerConfig, String> {
let value = get_config_value(iii).await?;
if value.is_null() {
tracing::info!("no configuration value found; using built-in default configuration");
Expand All @@ -70,22 +73,22 @@ pub async fn fetch_config(iii: &III) -> Result<WorkerConfig, String> {
WorkerConfig::from_json(&value)
}

async fn should_seed_default_value(iii: &III) -> Result<bool, String> {
async fn should_seed_default_value(iii: &IIIClient) -> Result<bool, String> {
match try_get_config_value(iii).await? {
None => Ok(true),
Some(value) if value.is_null() => Ok(true),
Some(_) => Ok(false),
}
}

async fn get_config_value(iii: &III) -> Result<Value, String> {
async fn get_config_value(iii: &IIIClient) -> Result<Value, String> {
try_get_config_value(iii)
.await?
.ok_or_else(|| format!("configuration `{CONFIG_ID}` not found"))
}

/// Returns `Ok(None)` when the entry does not exist (codes vary in case).
async fn try_get_config_value(iii: &III) -> Result<Option<Value>, String> {
async fn try_get_config_value(iii: &IIIClient) -> Result<Option<Value>, String> {
match trigger_with_retry(iii, "configuration::get", json!({ "id": CONFIG_ID })).await {
Ok(resp) => Ok(resp.get("value").cloned()),
Err(e) if e.to_ascii_uppercase().contains("NOT_FOUND") => Ok(None),
Expand Down Expand Up @@ -121,7 +124,7 @@ fn legacy_harness_keys(value: &Value) -> Vec<&'static str> {
/// can succeed. There is no `configuration::delete`; unlock a permissive schema,
/// replace the value with built-in defaults, then let the caller register the
/// real schema.
async fn discard_legacy_harness_config_if_needed(iii: &III) -> Result<(), String> {
async fn discard_legacy_harness_config_if_needed(iii: &IIIClient) -> Result<(), String> {
let Some(value) = try_get_config_value(iii).await? else {
return Ok(());
};
Expand Down Expand Up @@ -178,7 +181,7 @@ pub struct TriggerHandles {
/// Best-effort binding: the cron trigger type always exists (engine built-in),
/// but a transient failure must not brick boot — it surfaces as a `None`
/// handle.
fn bind(iii: &III, trigger_type: &str, function_id: &str, config: Value) -> Option<Trigger> {
fn bind(iii: &IIIClient, trigger_type: &str, function_id: &str, config: Value) -> Option<Trigger> {
match iii.register_trigger(RegisterTriggerInput {
trigger_type: trigger_type.to_string(),
function_id: function_id.to_string(),
Expand All @@ -197,7 +200,7 @@ fn bind(iii: &III, trigger_type: &str, function_id: &str, config: Value) -> Opti
}

/// (Re)bind the cron pending-sweep from the current config.
pub fn bind_sweep(iii: &III, cfg: &WorkerConfig) -> Option<Trigger> {
pub fn bind_sweep(iii: &IIIClient, cfg: &WorkerConfig) -> Option<Trigger> {
bind(
iii,
"cron",
Expand Down Expand Up @@ -238,10 +241,10 @@ pub struct OnConfigChangeResponse {
/// trigger. `handles` holds the live cron `Trigger` the handler re-binds when
/// `sweep_expression` changes.
pub fn register_config_trigger(
iii: &Arc<III>,
iii: &Arc<IIIClient>,
cell: ConfigCell,
handles: Arc<TriggerHandles>,
) -> Result<(), IIIError> {
) -> Result<(), Error> {
let cell_for_fn = cell.clone();
let handles_for_fn = handles.clone();
let engine = iii.clone();
Expand All @@ -253,7 +256,7 @@ pub fn register_config_trigger(
let engine = engine.clone();
async move {
on_config_change(&engine, &cell, &handles).await;
Ok::<OnConfigChangeResponse, IIIError>(OnConfigChangeResponse { ok: true })
Ok::<OnConfigChangeResponse, Error>(OnConfigChangeResponse { ok: true })
}
})
.description(
Expand All @@ -277,7 +280,7 @@ pub fn register_config_trigger(

/// Reload from the AUTHORITATIVE configuration. The caller-supplied trigger
/// payload is intentionally ignored: a direct call can never inject config.
async fn on_config_change(iii: &III, cell: &ConfigCell, handles: &TriggerHandles) {
async fn on_config_change(iii: &IIIClient, cell: &ConfigCell, handles: &TriggerHandles) {
let cfg = match fetch_config(iii).await {
Ok(cfg) => cfg,
Err(e) => {
Expand All @@ -296,7 +299,11 @@ async fn on_config_change(iii: &III, cell: &ConfigCell, handles: &TriggerHandles
tracing::info!("harness configuration reloaded");
}

async fn trigger_with_retry(iii: &III, function_id: &str, payload: Value) -> Result<Value, String> {
async fn trigger_with_retry(
iii: &IIIClient,
function_id: &str,
payload: Value,
) -> Result<Value, String> {
let mut last_err = String::new();
for attempt in 1..=CONFIG_RETRIES {
match iii
Expand Down
6 changes: 3 additions & 3 deletions harness/src/deps.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@

use std::sync::Arc;

use iii_sdk::III;
use iii_sdk::IIIClient;

use crate::clients::{
ContextClient, EngineClient, FunctionDescriptor, RouterClient, SessionClient,
Expand All @@ -18,7 +18,7 @@ use crate::locks::SessionLocks;

#[derive(Clone)]
pub struct Deps {
pub iii: Arc<III>,
pub iii: Arc<IIIClient>,
pub config: ConfigCell,
pub functions: FunctionsCell,
pub events: TurnEvents,
Expand All @@ -28,7 +28,7 @@ pub struct Deps {

impl Deps {
pub fn new(
iii: Arc<III>,
iii: Arc<IIIClient>,
config: ConfigCell,
functions: FunctionsCell,
events: TurnEvents,
Expand Down
12 changes: 7 additions & 5 deletions harness/src/discovery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@

use std::sync::Arc;

use iii_sdk::{IIIError, RegisterFunction, RegisterTriggerInput, III};
use iii_sdk::errors::Error;
use iii_sdk::protocol::RegisterTriggerInput;
use iii_sdk::{IIIClient, RegisterFunction};
use serde_json::json;
use tokio::sync::RwLock;

Expand All @@ -40,7 +42,7 @@ pub async fn apply(cell: &FunctionsCell, functions: Vec<FunctionDescriptor>) {
}

/// Fetch the authoritative registry and swap the snapshot; returns the count.
async fn reload(iii: &Arc<III>, cell: &FunctionsCell, timeout_ms: u64) -> usize {
async fn reload(iii: &Arc<IIIClient>, cell: &FunctionsCell, timeout_ms: u64) -> usize {
let engine = EngineClient::new(iii.clone(), timeout_ms);
let functions = engine.functions_list().await;
let count = functions.len();
Expand All @@ -50,7 +52,7 @@ async fn reload(iii: &Arc<III>, cell: &FunctionsCell, timeout_ms: u64) -> usize

/// Seed the snapshot from the registry. The trigger fires only on change, so
/// without this the cache would stay empty until the first change.
pub async fn seed(iii: &Arc<III>, cell: &FunctionsCell, timeout_ms: u64) {
pub async fn seed(iii: &Arc<IIIClient>, cell: &FunctionsCell, timeout_ms: u64) {
let count = reload(iii, cell, timeout_ms).await;
tracing::info!(count, "seeded function-registry cache");
}
Expand All @@ -75,7 +77,7 @@ pub struct OnFunctionsChangeResponse {
/// seeded snapshot still serves — it just won't update until restart) rather
/// than bricking boot. The handler is tagged `internal` so it stays off the
/// public catalog (and out of the very cache it maintains).
pub fn register_functions_trigger(iii: &Arc<III>, cell: FunctionsCell, timeout_ms: u64) {
pub fn register_functions_trigger(iii: &Arc<IIIClient>, cell: FunctionsCell, timeout_ms: u64) {
let engine = iii.clone();
iii.register_function(
FUNCTIONS_FN_ID,
Expand All @@ -85,7 +87,7 @@ pub fn register_functions_trigger(iii: &Arc<III>, cell: FunctionsCell, timeout_m
async move {
let count = reload(&engine, &cell, timeout_ms).await;
tracing::debug!(count, "function-registry cache refreshed");
Ok::<OnFunctionsChangeResponse, IIIError>(OnFunctionsChangeResponse { ok: true })
Ok::<OnFunctionsChangeResponse, Error>(OnFunctionsChangeResponse { ok: true })
}
})
.description(
Expand Down
Loading
Loading