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
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ Multi-crate Rust workspace for the PromptForge pipeline engine, the harness that
- PromptForge crates are named `promptforge` and promptforge-* and must not depend on gateway, workshop, or harness crates
- PromptForge has one public crate, `promptforge` at crates/promptforge/: a facade of single-item re-exports grouped into documented role modules. Crates outside the family may depend only on `promptforge`, never on a promptforge-* crate. Everything else lives under crates/promptforge-internal/, a manifestless container private to the family that holds the engine (`promptforge-engine`), the types crate (`promptforge-types`), the virtual filesystem (`promptforge-vfs`), and the lua, parser, and model-client crates; `promptforge` is the only outside crate permitted to depend into it
- The desktop app (the `workshop` crate) depends on `workshop-server-api` and never on `workshop-server`; the facade is the desktop app's entire view of the server
- Shared crates are named shared-*, contain the public API surface across products and downstream crates, and must not depend on any product crates. PromptForge's own public surface is the `promptforge` facade, and its types crate (`promptforge-types`) has left shared-* for the private container; Gateway's is gateway-api-types and gateway-api-discovery, named gateway-* now that both have left shared-*; the types crate contains the wire vocabulary only, never code
- Shared crates are named shared-*, contain the public API surface across products and downstream crates, and must not depend on any product crates. PromptForge's own public surface is the `promptforge` facade, and its types crate (`promptforge-types`) has left shared-* for the private container; Gateway's is gateway-api-types and gateway-api-discovery, named gateway-* now that both have left shared-*; `promptforge-types` holds the shared vocabulary together with host-support code (the untrusted guards, the cancellation tree, and the event emitter)
- Crates named build-* are for building specific outputs
- Dependency rules bind all kinds: normal, dev, build, and target-specific dependencies. One exception: a crate under crates/promptforge-internal/ may list `promptforge` in `[dev-dependencies]` only so its doc examples compile against the facade paths hosts see. No unit test, integration test, or bench imports it. This edge is exempt from the one-way flow rule under Structural Rules.

Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

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

1 change: 1 addition & 0 deletions crates/harness-internal/models/tests/it/end_to_end.rs
Original file line number Diff line number Diff line change
Expand Up @@ -322,6 +322,7 @@ async fn a_prepared_run_drives_end_to_end_and_records_the_whole_stream() {
let services = Services {
registry: None,
vfs: promptforge::vfs::VfsRef::default(),
input_text: None,
cancel: CancelHandle::new(),
log: Arc::clone(&log),
chat: Arc::new(GatewayChatPerformer::new(client, deltas)),
Expand Down
99 changes: 99 additions & 0 deletions crates/harness-internal/runner/src/files.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
//! The host's half of a prompt's declared store files: the launch's input
//! text staged at the frontmatter's `input:` path before a run, and the
//! `output:` path read after it.
//!
//! Both go through the handle's store view ([`VfsRef::acquire_store`]),
//! so the store's strict path rules apply to the declared paths, which
//! the parser takes as written, and the handle's policy and op sink see
//! each operation like any other. Both are synchronous, like the VFS; a
//! caller runs them on the blocking pool.

use promptforge::vfs::{Origin, StoreOp, StoreOutcome, VfsError, VfsRef, perform_store_op};

/// Why a run's declared input file could not be put in place.
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum InputFileError {
/// The launch supplied input text, but the prompt declares no
/// `input:` file to stage it at.
#[error("the launch supplied input text, but the prompt declares no `input:` file")]
Undeclared,
/// The prompt declares an input file that the launch did not supply
/// and the store does not already hold.
#[error(
"the prompt declares the input file `{path}`, but the launch supplied no input text \
and the store has no such file"
)]
Missing {
/// The declared input path.
path: String,
},
/// The store refused the staging write or the existence check.
#[error("the input file `{path}` could not be staged")]
Store {
/// The declared input path.
path: String,
/// The store's failure.
#[source]
source: VfsError,
},
}

/// Puts the prompt's declared input in place in `vfs`'s store: writes
/// `text` at `declared`, overwriting any file there, or, when the launch
/// supplied no text, checks that the store already holds the file.
///
/// # Errors
/// Returns [`InputFileError::Undeclared`] for text the prompt declares no
/// file for, [`InputFileError::Missing`] for a declared file that is
/// neither supplied nor present, and [`InputFileError::Store`] when the
/// store refuses the write or the check.
pub fn stage_input(
vfs: &VfsRef,
declared: Option<&str>,
text: Option<String>,
) -> Result<(), InputFileError> {
let path = match (declared, &text) {
(None, None) => return Ok(()),
(None, Some(_)) => return Err(InputFileError::Undeclared),
(Some(path), _) => path,
};
let store = |source| InputFileError::Store {
path: path.to_owned(),
source,
};
let view = vfs
.acquire_store(Origin::new(format!("input: {path}")))
.map_err(store)?;
let path = path.to_owned();
match text {
Some(contents) => perform_store_op(&view, StoreOp::Write { path, contents })
.map(drop)
.map_err(store),
None => match perform_store_op(&view, StoreOp::Exists { path: path.clone() }) {
Ok(StoreOutcome::Bool(true)) => Ok(()),
Ok(_) => Err(InputFileError::Missing { path }),
Err(source) => Err(store(source)),
},
}
}

/// Reads the prompt's declared output file at `path` from `vfs`'s store.
///
/// # Errors
/// Returns the store's failure, [`VfsError::NotFound`] when the run never
/// wrote the file.
pub fn read_output(vfs: &VfsRef, path: &str) -> Result<String, VfsError> {
let view = vfs.acquire_store(Origin::new(format!("output: {path}")))?;
let read = StoreOp::Read {
path: path.to_owned(),
start: None,
end: None,
};
match perform_store_op(&view, read)? {
StoreOutcome::Text(text) => Ok(text),
other => Err(VfsError::Backend {
message: format!("a whole-file store read answered {other:?} instead of text"),
}),
}
}
4 changes: 3 additions & 1 deletion crates/harness-internal/runner/src/lib.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
//! harness-runner - the harness effect loop: prepares an engine `Run`
//! from a prompt file (drawing the host inputs the engine refuses to draw
//! itself, activating capabilities, opening the run's row), steps it,
//! itself, putting the declared input file in place, activating
//! capabilities, opening the run's row), steps it,
//! performs each effect on tokio through one performer per effect kind,
//! feeds the answers back, records every event, effect, and answer in the
//! run log, and owns cancellation.
Expand Down Expand Up @@ -35,6 +36,7 @@
pub mod cancel;
mod display_chain;
pub mod effect_loop;
pub mod files;
pub mod performers;
pub mod prepare;
pub mod spawn;
Expand Down
112 changes: 96 additions & 16 deletions crates/harness-internal/runner/src/prepare.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,9 @@
//! CSPRNG and its `started_at` from the wall clock, both written to the
//! run's row in the log before anything else, so the record can hand them
//! back verbatim to a future replay. Then the ceremony the engine's
//! `Environment` expects of a host: parse; hand the run's whole
//! filesystem, host roots and the declared store, to the capabilities'
//! `Environment` expects of a host: parse; put the prompt's declared
//! `input:` file in place in the store (`files::stage_input`); hand the
//! run's whole filesystem, host roots and the declared store, to the capabilities'
//! services and to the context as given; activate the prompt's declared
//! capabilities against the caller's registry, which assembles the
//! catalog, the preludes, and the implementation table; install the
Expand All @@ -16,7 +17,8 @@
//! unsatisfiable prompt with the engine's model-readable notice; and
//! build the `Run` beside its performers.
//!
//! A refusal (or a prompt that fails to parse) is a run that ended before
//! A refusal (or a prompt that fails to parse, or an input file that
//! cannot be put in place) is a run that ended before
//! it began: its row is closed as failed with the refusal as the message,
//! so the log answers "why did this session fail" for a run the loop
//! never saw.
Expand All @@ -33,16 +35,18 @@ use promptforge::cancel::CancelHandle;
use promptforge::event::Event;
use promptforge::model::ModelDescriptor;
use promptforge::timestamp::Timestamp;
use promptforge::vfs::VfsRef;
use promptforge::vfs::{VfsError, VfsRef};
use promptforge::{Environment, RunContext, RunError};
use promptforge::{ParseError, Prompt, Run};
use sha2::{Digest as _, Sha256};

use crate::display_chain::display_chain;
use crate::effect_loop::{SharedLog, failed_outcome};
use crate::files::{InputFileError, stage_input};
use crate::performers::{
ActivatedTools, ChatPerformer, LogTaskEvents, Performers, TokioTimer, VfsStore,
};
use crate::spawn::spawn_blocking_launch;

/// What the caller owns and preparation borrows: the registry of
/// installed capabilities, the host roots, the run's cancel flag, the
Expand All @@ -57,6 +61,9 @@ pub struct Services {
/// passed straight to the context's VFS and handed to the capabilities
/// as the run's services.
pub vfs: VfsRef,
/// The text staged at the prompt's declared `input:` path before the
/// run, when the launch supplied one.
pub input_text: Option<String>,
/// The run's cancel flag: handed to the context, to every capability
/// activated for the run, and polled by the engine.
pub cancel: CancelHandle,
Expand Down Expand Up @@ -112,6 +119,9 @@ pub struct Prepared {
/// run's own events; the caller hands them to its sink so the session
/// sees them in order.
pub parse_events: Vec<Event>,
/// The prompt's declared `output:` path, which the caller reads with
/// [`read_output`](crate::files::read_output) once the run completes.
pub output_path: Option<String>,
}

/// Why a run could not be prepared.
Expand Down Expand Up @@ -153,6 +163,18 @@ pub enum PrepareError {
#[source]
error: RunError,
},
/// The prompt's declared input file could not be put in place: the
/// launch supplied text the prompt declares no file for, the prompt
/// declares a file that is neither supplied nor in the store, or the
/// store refused. The run's row is closed as failed with kind `Input`.
#[error("the prompt's declared input file cannot be put in place")]
Input {
/// The run's row, closed with this refusal.
run_id: RunId,
/// Why the input could not be put in place.
#[source]
source: InputFileError,
},
/// The run log refused a write; the run cannot be recorded, so it is
/// not prepared.
#[error(transparent)]
Expand All @@ -161,16 +183,18 @@ pub enum PrepareError {

/// Prepares the prompt at `prompt_path` for one run with `args`: draws the
/// run's seed and start and opens its row in the log, parses the prompt,
/// activates its declared capabilities against the caller's registry,
/// prepares the context, refuses an unsatisfiable prompt, and builds the
/// `Run` and its performers.
/// puts its declared input file in place, activates its declared
/// capabilities against the caller's registry, prepares the context,
/// refuses an unsatisfiable prompt, and builds the `Run` and its
/// performers.
///
/// # Errors
/// Returns [`PrepareError::Read`] when the file cannot be read (no row is
/// written), [`PrepareError::Parse`] when it does not parse and
/// [`PrepareError::Refused`] when the environment cannot satisfy it (in
/// both cases the row is closed as failed), and [`PrepareError::Log`]
/// when the log refuses a write.
/// written), [`PrepareError::Parse`] when it does not parse,
/// [`PrepareError::Input`] when its declared input file cannot be put in
/// place, and [`PrepareError::Refused`] when the environment cannot
/// satisfy it (in these three cases the row is closed as failed), and
/// [`PrepareError::Log`] when the log refuses a write.
pub async fn prepare_run(
prompt_path: &Path,
args: &str,
Expand All @@ -191,10 +215,12 @@ pub async fn prepare_run(
/// source is attributed to in [`PrepareError::Parse`].
///
/// # Errors
/// Returns [`PrepareError::Parse`] when the source does not parse and
/// [`PrepareError::Refused`] when the environment cannot satisfy it (in
/// both cases the row is closed as failed), and [`PrepareError::Log`]
/// when the log refuses a write. Never [`PrepareError::Read`].
/// Returns [`PrepareError::Parse`] when the source does not parse,
/// [`PrepareError::Input`] when its declared input file cannot be put in
/// place, and [`PrepareError::Refused`] when the environment cannot
/// satisfy it (in these three cases the row is closed as failed), and
/// [`PrepareError::Log`] when the log refuses a write. Never
/// [`PrepareError::Read`].
pub async fn prepare_source(
source: &str,
prompt_path: &Path,
Expand All @@ -204,6 +230,7 @@ pub async fn prepare_source(
let Services {
registry,
vfs,
input_text,
cancel,
log,
chat,
Expand All @@ -223,7 +250,7 @@ pub async fn prepare_source(
.await
.begin_run(RunMeta {
session_id: session_id.clone(),
agent,
agent: agent.clone(),
prompt_hash: prompt_hash(source),
seed,
flags: 0,
Expand Down Expand Up @@ -256,6 +283,15 @@ pub async fn prepare_source(
}
};

// The declared input is in place before anything else sees the
// store: the capabilities activate over the same filesystem, and the
// run's first section may read it.
stage_declared_input(&prompt, &vfs, input_text, &agent, &log, run_id).await?;
let output_path = prompt
.frontmatter()
.output()
.map(|decl| decl.path().to_owned());

// The parse events were stamped under task `0` from zero; the run's
// root task continues the sequence past them, so `(task_id, task_seq)`
// is unique across every record of the run.
Expand Down Expand Up @@ -306,7 +342,51 @@ pub async fn prepare_source(
started_at,
performers,
parse_events,
output_path,
})
}

/// Puts `prompt`'s declared input file in place in `vfs`'s store on the
/// blocking pool, tagged with `agent`. A refusal closes `run_id`'s row as
/// failed under the `Input` kind, with the cause chain as its message.
async fn stage_declared_input(
prompt: &Prompt,
vfs: &VfsRef,
input_text: Option<String>,
agent: &str,
log: &SharedLog,
run_id: RunId,
) -> Result<(), PrepareError> {
let declared = prompt
.frontmatter()
.input()
.map(|decl| decl.path().to_owned());
if declared.is_none() && input_text.is_none() {
return Ok(());
}
let staging = vfs.clone();
let path = declared.clone();
let staged = spawn_blocking_launch(agent, move || {
stage_input(&staging, path.as_deref(), input_text)
})
.await
.unwrap_or_else(|join| {
Err(InputFileError::Store {
path: declared.unwrap_or_default(),
source: VfsError::Backend {
message: format!("the staging task failed: {join}"),
},
})
});
let Err(source) = staged else {
return Ok(());
};
let outcome = RunOutcome::Failed {
kind: "Input".to_owned(),
message: display_chain(&source),
};
close_failed(log, run_id, outcome).await?;
Err(PrepareError::Input { run_id, source })
}

/// Closes `run_id`'s row with `outcome`, a run that ended before the loop
Expand Down
9 changes: 6 additions & 3 deletions crates/harness-internal/runner/src/spawn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
//! [`Provenance`] - so a run's tasks trace as a group and slice by task.
//! The last two cover the work that performs no effect: a session's
//! supervisor, whose span records the session id, and a launch's
//! filesystem probes, whose span records the agent name. Each is a
//! filesystem work, whose span records the agent name. Each is a
//! permitted caller of the raw tokio method it wraps, and no other
//! harness code is.

Expand Down Expand Up @@ -110,8 +110,11 @@ where
/// a span named `launch` that records the agent name under `agent`.
///
/// A launch walks the agents directory and reads the agent's source
/// before any run or session exists, so the work has no [`Tag`] and no
/// session id; the agent name is what ties it to the launch that asked.
/// before any run or session exists, and each run puts the prompt's
/// declared input file in place before its first step and reads its
/// declared output file after its last. None of that performs an effect,
/// so the work has no [`Tag`]; the agent name is what ties it to the
/// launch that asked.
/// The closure runs to completion even if its [`JoinHandle`] is aborted
/// or dropped, just as with `tokio::task::spawn_blocking`.
///
Expand Down
Loading
Loading