use super::error::PluginHostImportErrorKind;
use super::{
PluginBufferObservationError, PluginHostImportBatches, PluginHostImportDiscardReport,
PluginHostImportError, PluginHostImportSession, PluginHostImportSnapshot,
PluginHostOwnerBatches, PluginHostOwnerDiscardReport, PluginHostStateMismatch,
};
use crate::plugin::host::handles::PluginBufferSnapshot;
use crate::plugin::host::workspace_io::PendingPluginWorkspaceObserveTaskBatch;
use crate::{
StatusInfoText,
buffer::BufferEdit,
ecs::{
components::buffer::{BufferEntity, ViewEntity},
events::status::{StatusMessage, StatusMessageRequested},
},
plugin::{
GuestDecodeError, GuestDecodeField, PendingPluginUpdateBatch,
PendingPluginWorkspaceIoBatch, PluginAuthorizationError, PluginCapabilities,
PluginCapabilityRef, PluginCapabilityShape, PluginEffectBatchLimit,
PluginEffectBatchLimitField, PluginEffectCommitError, PluginHandleError, PluginHandleStore,
PluginHostContext, PluginHostResourceKind, PluginIdentity,
PluginOperationalImportRejection, PluginResourceHandleRejectionReason,
PluginWorkspaceIoBudget, PluginWorkspaceIoBudgetField, PluginWorkspaceIoCommitError,
PluginWorkspaceObserveOutcome, PluginWorkspaceObserveTaskPoll,
PluginWorkspaceObserveTaskQueue, PluginWorkspaceObserveTaskQueueError,
PluginWorkspaceObserveTaskReserveError, WorkspacePathError, WorkspacePathGrant,
WorkspacePathRef,
},
text_stream::TextRevision,
};
use bevy::prelude::Entity;
use std::{
any::Any,
fs,
num::{NonZeroU64, NonZeroUsize},
panic::{self, AssertUnwindSafe},
path::{Path, PathBuf},
time::{SystemTime, UNIX_EPOCH},
};
fn assert_send_sync<T: Send + Sync>() {}
fn count_limit(value: usize) -> NonZeroUsize {
NonZeroUsize::new(value).expect("test limit should be non-zero")
}
fn assert_import_batch_debug_redacts_nested_work(debug: &str) {
for hidden in [
"SealedPluginEffectBatch",
"SealedPluginWorkspaceIoBatch",
"SealedPluginWorkspaceObserveTaskBatch",
"effect_shapes",
"request_shapes",
"secret debug status",
"secret-debug-output",
"secret debug bytes",
"secret-debug-input",
] {
assert!(!debug.contains(hidden), "debug output leaked {hidden:?}");
}
}
#[test]
fn host_import_batches_are_send_sync() {
assert_send_sync::<PluginHostImportBatches>();
assert_send_sync::<PluginHostOwnerBatches>();
}
#[test]
fn paired_host_import_proofs_enforce_one_identity() {
assert_panics_with(
|| {
let _batches = PluginHostImportBatches::new(
empty_effect_batch("left"),
empty_workspace_io_batch("right"),
empty_workspace_observe_task_batch("left"),
);
},
"plugin host import batches must share one identity",
);
assert_panics_with(
|| {
let _batches = PluginHostImportBatches::new(
empty_effect_batch("left"),
empty_workspace_io_batch("left"),
empty_workspace_observe_task_batch("right"),
);
},
"plugin host import batches must share one identity",
);
assert_panics_with(
|| {
let _batches = PluginHostOwnerBatches::from_parts(
empty_effect_batch("left"),
empty_workspace_io_batch("right"),
);
},
"plugin owner batches must share one identity",
);
assert_panics_with(
|| {
let _report = PluginHostOwnerDiscardReport::new(
empty_effect_batch("left").discard(),
empty_workspace_io_batch("right").discard(),
);
},
"plugin owner discard report must share one identity",
);
assert_panics_with(
|| {
let _report = PluginHostImportDiscardReport::new(
empty_effect_batch("left").discard(),
empty_workspace_io_batch("left").discard(),
empty_workspace_observe_task_batch("right").discard(),
);
},
"plugin host import discard report must share one identity",
);
}
#[test]
fn host_import_batches_split_owner_work_as_one_identity_proof() {
let view = ViewEntity(test_entity(2));
let host = PluginHostContext::for_test(
"owners",
PluginCapabilities {
status_publish: true,
workspace_observe: vec![WorkspacePathGrant::new("docs")],
workspace_artifact_write: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let mut handles = handle_store("owners");
let view_handle = handles.issue_view(view).expect("view handle");
let mut queue = PluginWorkspaceObserveTaskQueue::default();
let mut session = session(&host, &handles);
session
.push_status_text(view_handle, "owner status")
.expect("status import should queue");
session
.queue_workspace_artifact_write(
WorkspacePathRef::try_from("docs/output.txt").expect("path"),
b"owner bytes".to_vec(),
)
.expect("workspace write should queue");
let _task = session
.queue_workspace_observe_task(
WorkspacePathRef::try_from("docs/input.txt").expect("path"),
&mut queue,
)
.expect("task request should queue");
let update = session.seal().into_owner_and_workspace_observe_tasks();
assert_eq!(update.owner_batches.identity(), "owners");
assert_eq!(update.owner_batches.effects().identity(), "owners");
assert_eq!(update.owner_batches.workspace_io().identity(), "owners");
assert_eq!(update.workspace_observe_tasks.identity(), "owners");
assert_eq!(update.owner_batches.effects().effect_count(), 1);
assert_eq!(update.owner_batches.workspace_io().request_count(), 1);
assert_eq!(update.workspace_observe_tasks.task_count(), 1);
assert_eq!(queue.pending_len(), 0);
assert_eq!(queue.reserved_pending_len(), 1);
}
#[test]
fn host_import_batch_debug_output_is_count_shaped() {
let (batches, _queue) = host_import_batches_with_secret_work("debug-batches");
let debug = format!("{batches:?}");
assert!(debug.contains("PluginHostImportBatches"));
assert!(debug.contains("effect_count: 1"));
assert!(debug.contains("workspace_io_request_count: 1"));
assert!(debug.contains("workspace_observe_task_count: 1"));
assert_import_batch_debug_redacts_nested_work(&debug);
let update = batches.into_owner_and_workspace_observe_tasks();
let debug = format!("{:?}", update.owner_batches);
assert!(debug.contains("PluginHostOwnerBatches"));
assert!(debug.contains("effect_count: 1"));
assert!(debug.contains("workspace_io_request_count: 1"));
assert!(!debug.contains("workspace_observe_task_count"));
assert_import_batch_debug_redacts_nested_work(&debug);
}
#[test]
fn host_import_discard_debug_output_is_count_shaped() {
let (session, _queue) = host_import_session_with_secret_work("debug-discard");
let discard = session.discard();
assert_eq!(discard.identity(), "debug-discard");
assert_eq!(discard.effects().discarded_effects(), 1);
assert_eq!(discard.workspace_io().discarded_requests(), 1);
assert_eq!(discard.workspace_observe_tasks().discarded_tasks(), 1);
let debug = format!("{discard:?}");
assert!(debug.contains("PluginHostImportDiscardReport"));
assert!(debug.contains("discarded_effect_count: 1"));
assert!(debug.contains("discarded_workspace_io_request_count: 1"));
assert!(debug.contains("discarded_workspace_observe_task_count: 1"));
assert!(!debug.contains("PluginEffectDiscardReport"));
assert!(!debug.contains("PluginWorkspaceIoDiscardReport"));
assert!(!debug.contains("PluginWorkspaceObserveTaskDiscardReport"));
assert_import_batch_debug_redacts_nested_work(&debug);
let (batches, _queue) = host_import_batches_with_secret_work("debug-owner-discard");
let update = batches.into_owner_and_workspace_observe_tasks();
let owner_discard = update.owner_batches.discard();
let debug = format!("{owner_discard:?}");
assert!(debug.contains("PluginHostOwnerDiscardReport"));
assert!(debug.contains("discarded_effect_count: 1"));
assert!(debug.contains("discarded_workspace_io_request_count: 1"));
assert!(!debug.contains("discarded_workspace_observe_task_count"));
assert!(!debug.contains("PluginEffectDiscardReport"));
assert!(!debug.contains("PluginWorkspaceIoDiscardReport"));
assert_import_batch_debug_redacts_nested_work(&debug);
}
#[test]
fn status_import_authorizes_resolves_and_queues_effect() {
let view = ViewEntity(test_entity(2));
let host = PluginHostContext::for_test(
"status",
PluginCapabilities {
status_publish: true,
..PluginCapabilities::default()
},
);
let mut handles = handle_store("status");
let view_handle = handles.issue_view(view).expect("view handle");
let mut session = session(&host, &handles);
session
.push_status_text(view_handle, "plugin ok")
.expect("status import should queue");
let batches = session.seal().into_parts();
let report = batches.effects.drain();
assert_eq!(
report.status_messages(),
&[StatusMessageRequested {
target: view,
message: StatusMessage::Info(StatusInfoText::escaped_display("plugin ok")),
}]
);
assert!(report.buffer_edits().is_empty());
}
#[test]
fn denied_status_import_produces_no_effects() {
let view = ViewEntity(test_entity(2));
let host = PluginHostContext::for_test("status", PluginCapabilities::default());
let mut handles = handle_store("status");
let view_handle = handles.issue_view(view).expect("view handle");
let mut session = session(&host, &handles);
let error = session
.push_status_text(view_handle, "plugin ok")
.expect_err("status import should be denied");
let discard = session.discard();
assert_eq!(
error,
PluginHostImportError::Authorization(PluginAuthorizationError::Denied {
identity: PluginIdentity::try_new("status").expect("identity"),
capability: PluginCapabilityShape::StatusPublish,
})
);
assert_eq!(discard.effects().discarded_effects(), 0);
}
#[test]
fn import_snapshot_debug_redacts_handle_targets() {
let buffer = BufferEntity(test_entity(1));
let view = ViewEntity(test_entity(2));
let host = PluginHostContext::for_test(
"snapshot",
PluginCapabilities {
status_publish: true,
buffer_propose_edit: true,
..PluginCapabilities::default()
},
);
let mut handles = handle_store("snapshot");
let _buffer_handle = handles
.issue_buffer(buffer, Some(TextRevision::from(1)), Some(view))
.expect("buffer handle");
let _view_handle = handles.issue_view(view).expect("view handle");
let snapshot = PluginHostImportSnapshot::new(&host, &handles).expect("snapshot");
let debug = format!("{snapshot:?}");
assert!(debug.contains("PluginHostImportSnapshot"));
assert!(debug.contains("buffer_handles"));
assert!(!debug.contains("Entity"));
assert!(!debug.contains("BufferEntity"));
assert!(!debug.contains("ViewEntity"));
}
#[test]
fn import_snapshot_rejects_cross_plugin_host_state() {
let host = PluginHostContext::for_test(
"grants-owner",
PluginCapabilities {
status_publish: true,
..PluginCapabilities::default()
},
);
let handles = handle_store("handle-owner");
let error = PluginHostImportSnapshot::new(&host, &handles)
.expect_err("snapshot must not mix grants and handles across plugins");
assert_eq!(
error,
PluginHostStateMismatch::new(
PluginIdentity::try_new("grants-owner").expect("identity"),
PluginIdentity::try_new("handle-owner").expect("identity"),
)
);
assert_eq!(error.host_identity().as_str(), "grants-owner");
assert_eq!(error.handle_identity().as_str(), "handle-owner");
assert_eq!(
error.to_string(),
"host identity \"grants-owner\" does not match handle store identity \"handle-owner\""
);
}
#[test]
fn import_session_rejects_cross_plugin_host_state() {
let host = PluginHostContext::for_test("session-owner", PluginCapabilities::default());
let handles = handle_store("handle-owner");
let error = PluginHostImportSession::new(
&host,
&handles,
PluginEffectBatchLimit::default(),
PluginWorkspaceIoBudget::try_new(4, 64, 64).expect("budget"),
)
.expect_err("session must reject cross-plugin host state");
assert!(matches!(error, PluginHostImportError::HostStateMismatch(_)));
assert_eq!(
error.to_string(),
"plugin import host state denied: host identity \"session-owner\" does not match handle store identity \"handle-owner\""
);
}
#[test]
fn import_snapshot_reuses_host_authority_projection() {
let host = PluginHostContext::for_test(
"workspace",
PluginCapabilities {
workspace_observe: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let handles = handle_store("workspace");
let snapshot = PluginHostImportSnapshot::new(&host, &handles).expect("snapshot");
let access = snapshot
.authorize_workspace_read(
WorkspacePathRef::try_from("docs/input.txt").expect("path should validate"),
)
.expect("workspace read should be authorized");
let denied = snapshot
.authorize(PluginCapabilityRef::StatusPublish)
.expect_err("status should be denied");
assert_eq!(snapshot.identity_proof().as_str(), "workspace");
assert_eq!(access.workspace_relative_path(), "docs/input.txt");
assert_eq!(
denied,
PluginAuthorizationError::Denied {
identity: PluginIdentity::try_new("workspace").expect("identity"),
capability: PluginCapabilityShape::StatusPublish,
}
);
}
#[test]
fn buffer_observe_import_returns_bounded_snapshot() {
let buffer = BufferEntity(test_entity(1));
let host = PluginHostContext::for_test(
"observe",
PluginCapabilities {
buffer_observe: true,
..PluginCapabilities::default()
},
);
let mut handles = handle_store("observe");
let handle = handles
.issue_observed_buffer(
buffer,
PluginBufferSnapshot::new(TextRevision::from(3), "secret text"),
None,
)
.expect("buffer handle");
let session = session(&host, &handles);
let snapshot = session
.observe_buffer(handle, message_limit(64))
.expect("buffer observe should return snapshot");
assert_eq!(snapshot.revision(), TextRevision::from(3));
assert_eq!(snapshot.text(), "secret text");
assert!(!format!("{snapshot:?}").contains("secret"));
}
#[test]
fn buffer_observe_import_requires_capability() {
let buffer = BufferEntity(test_entity(1));
let host = PluginHostContext::for_test("observe", PluginCapabilities::default());
let mut handles = handle_store("observe");
let handle = handles
.issue_observed_buffer(
buffer,
PluginBufferSnapshot::new(TextRevision::from(3), "secret text"),
None,
)
.expect("buffer handle");
let session = session(&host, &handles);
let error = session
.observe_buffer(handle, message_limit(64))
.expect_err("buffer observe should be denied");
assert_eq!(
error,
PluginHostImportError::Authorization(PluginAuthorizationError::Denied {
identity: PluginIdentity::try_new("observe").expect("identity"),
capability: PluginCapabilityShape::BufferObserve,
})
);
assert!(!format!("{error:?}").contains("secret"));
assert!(!error.to_string().contains("secret"));
}
#[test]
fn buffer_observe_import_rejects_missing_or_oversized_snapshot() {
let buffer = BufferEntity(test_entity(1));
let host = PluginHostContext::for_test(
"observe",
PluginCapabilities {
buffer_observe: true,
..PluginCapabilities::default()
},
);
let mut handles = handle_store("observe");
let missing = handles
.issue_buffer(buffer, Some(TextRevision::from(3)), None)
.expect("buffer handle");
let oversized = handles
.issue_observed_buffer(
buffer,
PluginBufferSnapshot::new(TextRevision::from(3), "secret text"),
None,
)
.expect("buffer handle");
let session = session(&host, &handles);
let missing_error = session
.observe_buffer(missing, message_limit(64))
.expect_err("missing snapshot should reject");
let oversized_error = session
.observe_buffer(oversized, message_limit(4))
.expect_err("oversized snapshot should reject");
assert!(matches!(
missing_error,
PluginHostImportError::Observation(PluginBufferObservationError::MissingSnapshot { .. })
));
assert!(matches!(
oversized_error,
PluginHostImportError::Observation(PluginBufferObservationError::TextTooLarge {
byte_len: 11,
max: 4,
..
})
));
assert!(!format!("{missing_error:?}").contains("secret"));
assert!(!format!("{oversized_error:?}").contains("secret"));
assert_eq!(
missing_error.to_string(),
"plugin import observation denied: plugin buffer observation for buffer handle lacks captured text"
);
assert!(!oversized_error.to_string().contains("secret"));
}
#[test]
fn buffer_edit_import_requires_observed_revision() {
let buffer = BufferEntity(test_entity(1));
let host = PluginHostContext::for_test(
"edit",
PluginCapabilities {
buffer_propose_edit: true,
..PluginCapabilities::default()
},
);
let mut handles = handle_store("edit");
let handle = handles
.issue_buffer(buffer, None, None)
.expect("buffer handle");
let mut session = session(&host, &handles);
let error = session
.propose_buffer_edit(handle, BufferEdit::insert(0, "plugin"))
.expect_err("edit without revision should fail");
assert_eq!(
error,
PluginHostImportError::MissingObservedRevision {
handle: handle.shape(),
}
);
assert_eq!(
error.to_string(),
"plugin buffer edit for buffer handle lacks observed revision"
);
}
#[test]
fn resource_handle_rejections_have_stable_redacted_text() {
for (reason, expected) in [
(PluginResourceHandleRejectionReason::Missing, "missing"),
(PluginResourceHandleRejectionReason::WrongType, "wrong-type"),
(
PluginResourceHandleRejectionReason::Unavailable,
"unavailable",
),
] {
assert_eq!(reason.as_str(), expected);
assert_eq!(reason.to_string(), expected);
}
let error = PluginHostImportError::ResourceHandle {
kind: PluginHostResourceKind::ViewHandle,
reason: PluginResourceHandleRejectionReason::WrongType,
};
assert_eq!(
error.to_string(),
"plugin import resource view denied: wrong-type"
);
}
#[test]
fn host_import_errors_project_closed_operational_rejection_shapes() {
let mut handles = handle_store("classifier");
let handle = handles
.issue_buffer(BufferEntity(test_entity(1)), None, None)
.expect("buffer handle")
.shape();
let identity = PluginIdentity::try_new("classifier").expect("identity");
let cases = [
(
PluginHostImportError::HostStateMismatch(PluginHostStateMismatch::new(
identity.clone(),
PluginIdentity::try_new("other").expect("identity"),
)),
Some(PluginOperationalImportRejection::HostStateMismatch),
),
(
PluginHostImportError::Authorization(PluginAuthorizationError::Denied {
identity: identity.clone(),
capability: PluginCapabilityShape::StatusPublish,
}),
None,
),
(
PluginHostImportError::Handle(PluginHandleError::ZeroLimit),
Some(PluginOperationalImportRejection::Handle),
),
(
PluginHostImportError::Effect(PluginEffectCommitError::ZeroLimit {
field: PluginEffectBatchLimitField::MaxEffects,
}),
Some(PluginOperationalImportRejection::EffectQueue),
),
(
PluginHostImportError::WorkspaceIo(PluginWorkspaceIoCommitError::ZeroLimit {
field: PluginWorkspaceIoBudgetField::MaxRequests,
}),
Some(PluginOperationalImportRejection::WorkspaceIoQueue),
),
(
PluginHostImportError::WorkspaceObserveTask(
PluginWorkspaceObserveTaskReserveError::WorkspaceIo(
PluginWorkspaceIoCommitError::TooManyRequests {
identity: identity.clone(),
limit: count_limit(1),
},
),
),
Some(PluginOperationalImportRejection::WorkspaceIoQueue),
),
(
PluginHostImportError::WorkspaceObserveTask(
PluginWorkspaceObserveTaskReserveError::Queue(
PluginWorkspaceObserveTaskQueueError::TaskIdExhausted { identity },
),
),
Some(PluginOperationalImportRejection::WorkspaceObserveTaskQueue),
),
(
PluginHostImportError::Observation(PluginBufferObservationError::MissingSnapshot {
handle,
}),
Some(PluginOperationalImportRejection::Observation),
),
(
PluginHostImportError::Decode(GuestDecodeError::PayloadTooLarge {
field: GuestDecodeField::WorkspaceArtifactWriteBytes,
size: 9,
max: 8,
}),
Some(PluginOperationalImportRejection::Decode),
),
(
PluginHostImportError::ResourceHandle {
kind: PluginHostResourceKind::BufferHandle,
reason: PluginResourceHandleRejectionReason::Missing,
},
Some(PluginOperationalImportRejection::ResourceHandle),
),
(
PluginHostImportError::SessionMissing,
Some(PluginOperationalImportRejection::SessionMissing),
),
(
PluginHostImportError::SessionAlreadyActive,
Some(PluginOperationalImportRejection::SessionAlreadyActive),
),
(
PluginHostImportError::MissingObservedRevision { handle },
Some(PluginOperationalImportRejection::MissingObservedRevision),
),
];
for (error, expected) in cases {
assert_eq!(error.operational_rejection(), expected);
}
}
#[test]
fn host_import_error_kinds_render_stable_text() {
let cases = [
(
PluginHostImportErrorKind::HostStateMismatch,
"host-state-mismatch",
),
(PluginHostImportErrorKind::Authorization, "authorization"),
(PluginHostImportErrorKind::Handle, "handle"),
(PluginHostImportErrorKind::Effect, "effect"),
(PluginHostImportErrorKind::WorkspaceIo, "workspace-io"),
(
PluginHostImportErrorKind::WorkspaceObserveTask,
"workspace-observe-task",
),
(PluginHostImportErrorKind::Observation, "observation"),
(PluginHostImportErrorKind::Decode, "decode"),
(PluginHostImportErrorKind::ResourceHandle, "resource-handle"),
(PluginHostImportErrorKind::SessionMissing, "session-missing"),
(
PluginHostImportErrorKind::SessionAlreadyActive,
"session-already-active",
),
(
PluginHostImportErrorKind::MissingObservedRevision,
"missing-observed-revision",
),
];
for (kind, expected) in cases {
assert_eq!(kind.as_str(), expected);
assert_eq!(kind.to_string(), expected);
assert_eq!(format!("{kind:?}"), expected);
}
}
#[test]
fn host_import_error_debug_uses_redacted_closed_shapes() {
let mut handles = handle_store("debugger");
let handle = handles
.issue_buffer(BufferEntity(test_entity(1)), None, None)
.expect("buffer handle")
.shape();
let identity = PluginIdentity::try_new("debugger").expect("identity");
let errors = [
PluginHostImportError::Authorization(PluginAuthorizationError::Denied {
identity,
capability: PluginCapabilityShape::WorkspaceArtifactWrite,
}),
PluginHostImportError::Decode(GuestDecodeError::InvalidWorkspacePath {
field: GuestDecodeField::WorkspaceArtifactWritePath,
source: WorkspacePathError::DotComponent,
}),
PluginHostImportError::Observation(PluginBufferObservationError::TextTooLarge {
handle,
byte_len: 21,
max: 4,
}),
PluginHostImportError::ResourceHandle {
kind: PluginHostResourceKind::ViewHandle,
reason: PluginResourceHandleRejectionReason::WrongType,
},
PluginHostImportError::MissingObservedRevision { handle },
];
for error in errors {
let debug = format!("{error:?}");
assert!(debug.contains("PluginHostImportError"));
assert!(debug.contains("kind"));
assert!(debug.contains("operational_rejection"));
assert!(!debug.contains("secret"));
assert!(!debug.contains("docs/../secret.txt"));
assert!(!debug.contains("replacement payload"));
}
}
#[test]
fn buffer_edit_import_queues_revision_guarded_effect() {
let buffer = BufferEntity(test_entity(1));
let host = PluginHostContext::for_test(
"edit",
PluginCapabilities {
buffer_propose_edit: true,
..PluginCapabilities::default()
},
);
let mut handles = handle_store("edit");
let handle = handles
.issue_buffer(buffer, Some(TextRevision::from(7)), None)
.expect("buffer handle");
let mut session = session(&host, &handles);
session
.propose_buffer_edit(handle, BufferEdit::insert(0, "plugin"))
.expect("edit import should queue");
let batches = session.seal().into_parts();
let report = batches.effects.drain();
assert_eq!(report.buffer_edits().len(), 1);
assert_eq!(report.buffer_edits()[0].target, buffer);
assert_eq!(
report.buffer_edits()[0].provenance.base_revision(),
TextRevision::from(7)
);
assert_eq!(
report.buffer_edits()[0].provenance.source_identity(),
"edit"
);
assert_eq!(
report.buffer_edits()[0].provenance.capability().as_str(),
"buffer.propose_edit"
);
assert!(report.status_messages().is_empty());
}
#[test]
fn workspace_import_uses_path_grant_and_deferred_policy_token() {
let host = PluginHostContext::for_test(
"workspace",
PluginCapabilities {
workspace_artifact_write: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let handles = handle_store("workspace");
let mut session = session(&host, &handles);
session
.queue_workspace_artifact_write(
WorkspacePathRef::try_from("docs/output.txt").expect("path"),
b"written".to_vec(),
)
.expect("workspace write should queue");
let batches = session.seal().into_parts();
let root = temp_dir("workspace-import");
fs::create_dir(root.join("docs")).expect("docs directory");
let filesystem =
crate::fs_utils::FilesystemConfig::from_workspace_root(&root).expect("filesystem");
let report = batches.workspace_io.execute_synchronous(&filesystem);
assert_eq!(report.completions().len(), 1);
assert!(
report.completions()[0]
.as_write()
.expect("completion should be a write")
.outcome()
.is_ok()
);
assert_eq!(
fs::read_to_string(root.join("docs/output.txt")).expect("output"),
"written"
);
let _cleanup = fs::remove_dir_all(root);
}
#[test]
fn workspace_observe_import_returns_guest_outcome_without_queueing_io() {
let root = temp_dir("workspace-observe-import");
fs::create_dir(root.join("docs")).expect("docs directory");
fs::write(root.join("docs/input.txt"), b"secret workspace bytes").expect("input file");
let filesystem =
crate::fs_utils::FilesystemConfig::from_workspace_root(&root).expect("filesystem");
let host = PluginHostContext::for_test(
"workspace",
PluginCapabilities {
workspace_observe: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let handles = handle_store("workspace");
let mut session = session(&host, &handles);
let outcome = session
.observe_workspace(
WorkspacePathRef::try_from("docs/input.txt").expect("path"),
&filesystem,
)
.expect("workspace observe should return a guest outcome");
let batches = session.seal().into_parts();
assert_eq!(outcome.bytes(), Some(b"secret workspace bytes".as_slice()));
assert_eq!(batches.workspace_io.request_count(), 0);
let debug = format!("{outcome:?}");
assert!(!debug.contains("secret workspace bytes"));
assert!(!debug.contains("input.txt"));
let _cleanup = fs::remove_dir_all(root);
}
#[test]
fn workspace_observe_import_maps_filesystem_failures_to_closed_rejection() {
let root = temp_dir("workspace-observe-missing");
fs::create_dir(root.join("docs")).expect("docs directory");
let filesystem =
crate::fs_utils::FilesystemConfig::from_workspace_root(&root).expect("filesystem");
let host = PluginHostContext::for_test(
"workspace",
PluginCapabilities {
workspace_observe: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let handles = handle_store("workspace");
let mut session = session(&host, &handles);
let outcome = session
.observe_workspace(
WorkspacePathRef::try_from("docs/missing.txt").expect("path"),
&filesystem,
)
.expect("filesystem read failures should not fail the update");
assert_eq!(outcome, PluginWorkspaceObserveOutcome::Rejected);
assert!(!format!("{outcome:?}").contains("missing.txt"));
let _cleanup = fs::remove_dir_all(root);
}
#[test]
fn workspace_observe_import_requires_capability() {
let root = temp_dir("workspace-observe-denied");
fs::create_dir(root.join("docs")).expect("docs directory");
let filesystem =
crate::fs_utils::FilesystemConfig::from_workspace_root(&root).expect("filesystem");
let host = PluginHostContext::for_test("workspace", PluginCapabilities::default());
let handles = handle_store("workspace");
let mut session = session(&host, &handles);
let error = session
.observe_workspace(
WorkspacePathRef::try_from("docs/input.txt").expect("path"),
&filesystem,
)
.expect_err("missing capability should reject before filesystem policy");
assert_eq!(
error,
PluginHostImportError::Authorization(PluginAuthorizationError::Denied {
identity: PluginIdentity::try_new("workspace").expect("identity"),
capability: PluginCapabilityShape::WorkspaceObserve,
})
);
let _cleanup = fs::remove_dir_all(root);
}
#[test]
fn workspace_observe_import_charges_request_budget_before_filesystem_access() {
let root = temp_dir("workspace-observe-budget");
fs::create_dir(root.join("docs")).expect("docs directory");
fs::write(root.join("docs/input.txt"), b"first observe").expect("input file");
let filesystem =
crate::fs_utils::FilesystemConfig::from_workspace_root(&root).expect("filesystem");
let host = PluginHostContext::for_test(
"workspace",
PluginCapabilities {
workspace_observe: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let handles = handle_store("workspace");
let mut session = PluginHostImportSession::new(
&host,
&handles,
PluginEffectBatchLimit::default(),
PluginWorkspaceIoBudget::try_new(1, 64, 64).expect("budget"),
)
.expect("session");
let first = session
.observe_workspace(
WorkspacePathRef::try_from("docs/input.txt").expect("path"),
&filesystem,
)
.expect("first observe should spend the request budget");
let second = session
.observe_workspace(
WorkspacePathRef::try_from("docs/missing.txt").expect("path"),
&filesystem,
)
.expect_err("second observe should reject before filesystem policy");
let batches = session.seal().into_parts();
assert_eq!(first.bytes(), Some(b"first observe".as_slice()));
assert_eq!(
second,
PluginHostImportError::WorkspaceIo(PluginWorkspaceIoCommitError::TooManyRequests {
identity: PluginIdentity::try_new("workspace").expect("identity"),
limit: count_limit(1),
})
);
assert_eq!(batches.workspace_io.request_count(), 0);
let _cleanup = fs::remove_dir_all(root);
}
#[test]
fn workspace_observe_task_import_requires_capability_before_reserving_queue() {
let host = PluginHostContext::for_test("workspace", PluginCapabilities::default());
let handles = handle_store("workspace");
let mut session = session(&host, &handles);
let mut queue = PluginWorkspaceObserveTaskQueue::default();
let next_task_id = queue.next_id_for_test();
let error = session
.queue_workspace_observe_task(
WorkspacePathRef::try_from("docs/input.txt").expect("path"),
&mut queue,
)
.expect_err("missing capability should reject before queue reservation");
let discard = session.discard();
assert_eq!(
error,
PluginHostImportError::Authorization(PluginAuthorizationError::Denied {
identity: PluginIdentity::try_new("workspace").expect("identity"),
capability: PluginCapabilityShape::WorkspaceObserve,
})
);
assert_eq!(discard.workspace_observe_tasks().discarded_tasks(), 0);
assert_eq!(discard.workspace_io().discarded_requests(), 0);
assert_eq!(queue.pending_len(), 0);
assert_eq!(queue.reserved_pending_len(), 0);
assert_eq!(queue.next_id_for_test(), next_task_id);
}
#[test]
fn workspace_observe_task_requests_are_discarded_before_successful_return() {
let host = PluginHostContext::for_test(
"workspace",
PluginCapabilities {
workspace_observe: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let handles = handle_store("workspace");
let mut session = session(&host, &handles);
let mut queue = PluginWorkspaceObserveTaskQueue::default();
let _task_resource_authority = session
.queue_workspace_observe_task(
WorkspacePathRef::try_from("docs/input.txt").expect("path"),
&mut queue,
)
.expect("task request should queue");
let discard = session.discard();
assert_eq!(discard.workspace_observe_tasks().identity(), "workspace");
assert_eq!(
discard.workspace_observe_tasks().identity_proof().as_str(),
"workspace"
);
assert_eq!(discard.workspace_observe_tasks().discarded_tasks(), 1);
assert_eq!(discard.workspace_io().discarded_requests(), 0);
assert_eq!(queue.pending_len(), 0);
assert_eq!(queue.reserved_pending_len(), 0);
}
#[test]
fn workspace_observe_task_requests_commit_only_after_successful_return() {
let root = temp_dir("workspace-observe-task-import");
fs::create_dir(root.join("docs")).expect("docs directory");
fs::write(root.join("docs/input.txt"), b"task bytes").expect("input file");
let filesystem =
crate::fs_utils::FilesystemConfig::from_workspace_root(&root).expect("filesystem");
let host = PluginHostContext::for_test(
"workspace",
PluginCapabilities {
workspace_observe: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let handles = handle_store("workspace");
let mut session = session(&host, &handles);
let mut queue = PluginWorkspaceObserveTaskQueue::default();
let task_resource_authority = session
.queue_workspace_observe_task(
WorkspacePathRef::try_from("docs/input.txt").expect("path"),
&mut queue,
)
.expect("task request should queue");
let authority_debug = format!("{task_resource_authority:?}");
assert!(authority_debug.contains("PluginWorkspaceObserveTaskResourceAuthority"));
assert!(!authority_debug.contains("docs/input.txt"));
let reserved_task_handle = task_resource_authority.handle().clone();
let reserved_task_id = reserved_task_handle.task_id();
assert_eq!(reserved_task_handle.identity_proof().as_str(), "workspace");
let batches = session.seal().into_parts();
assert_eq!(batches.workspace_io.request_count(), 0);
assert_eq!(batches.workspace_observe_tasks.identity(), "workspace");
assert_eq!(batches.workspace_observe_tasks.task_count(), 1);
let enqueue = batches
.workspace_observe_tasks
.enqueue(&mut queue)
.expect("sealed task should enqueue");
let task_id = enqueue.task_handles()[0].task_id();
assert_eq!(task_id, reserved_task_id);
assert_eq!(
enqueue.task_handles()[0].identity_proof().as_str(),
"workspace"
);
assert_eq!(queue.reserved_pending_len(), 0);
assert_eq!(enqueue.identity_proof().as_str(), "workspace");
assert_eq!(
queue.poll(&reserved_task_handle),
PluginWorkspaceObserveTaskPoll::Pending
);
let completion = queue
.execute_next(&filesystem)
.expect("task should execute after commit");
assert_eq!(completion.task_id(), task_id);
assert_eq!(completion.outcome().bytes(), Some(b"task bytes".as_slice()));
let _cleanup = fs::remove_dir_all(root);
}
fn session(host: &PluginHostContext, handles: &PluginHandleStore) -> PluginHostImportSession {
PluginHostImportSession::new(
host,
handles,
PluginEffectBatchLimit::default(),
PluginWorkspaceIoBudget::try_new(4, 64, 64).expect("budget"),
)
.expect("session")
}
fn host_import_batches_with_secret_work(
identity: &str,
) -> (PluginHostImportBatches, PluginWorkspaceObserveTaskQueue) {
let (session, queue) = host_import_session_with_secret_work(identity);
(session.seal(), queue)
}
fn host_import_session_with_secret_work(
identity: &str,
) -> (PluginHostImportSession, PluginWorkspaceObserveTaskQueue) {
let view = ViewEntity(test_entity(300));
let host = PluginHostContext::for_test(
identity,
PluginCapabilities {
status_publish: true,
workspace_observe: vec![WorkspacePathGrant::new("docs")],
workspace_artifact_write: vec![WorkspacePathGrant::new("docs")],
..PluginCapabilities::default()
},
);
let mut handles = handle_store(identity);
let view_handle = handles.issue_view(view).expect("view handle");
let mut queue = PluginWorkspaceObserveTaskQueue::default();
let mut session = session(&host, &handles);
session
.push_status_text(view_handle, "secret debug status")
.expect("status import should queue");
session
.queue_workspace_artifact_write(
WorkspacePathRef::try_from("docs/secret-debug-output.txt").expect("path"),
b"secret debug bytes".to_vec(),
)
.expect("workspace write should queue");
let _task = session
.queue_workspace_observe_task(
WorkspacePathRef::try_from("docs/secret-debug-input.txt").expect("path"),
&mut queue,
)
.expect("task request should queue");
(session, queue)
}
fn empty_effect_batch(identity: &str) -> crate::plugin::SealedPluginEffectBatch {
PendingPluginUpdateBatch::new(test_identity(identity), PluginEffectBatchLimit::default()).seal()
}
fn empty_workspace_io_batch(identity: &str) -> crate::plugin::SealedPluginWorkspaceIoBatch {
PendingPluginWorkspaceIoBatch::new(
test_identity(identity),
PluginWorkspaceIoBudget::try_new(4, 64, 64).expect("budget"),
)
.seal()
}
fn empty_workspace_observe_task_batch(
identity: &str,
) -> crate::plugin::SealedPluginWorkspaceObserveTaskBatch {
PendingPluginWorkspaceObserveTaskBatch::new(test_identity(identity)).seal()
}
fn handle_store(identity: &str) -> PluginHandleStore {
PluginHandleStore::new(test_identity(identity))
}
fn test_identity(identity: &str) -> PluginIdentity {
PluginIdentity::try_new(identity).expect("identity")
}
fn assert_panics_with(action: impl FnOnce(), expected: &str) {
let panic = panic::catch_unwind(AssertUnwindSafe(action))
.expect_err("expected invariant assertion to panic");
let message = panic_message(panic.as_ref());
assert!(
message.contains(expected),
"panic message {message:?} did not contain {expected:?}",
);
}
fn test_entity(index: u32) -> Entity {
Entity::from_raw_u32(index).expect("entity")
}
fn panic_message(payload: &(dyn Any + Send)) -> &str {
payload
.downcast_ref::<&str>()
.copied()
.or_else(|| payload.downcast_ref::<String>().map(String::as_str))
.unwrap_or("<non-string panic>")
}
fn message_limit(max: u64) -> NonZeroU64 {
NonZeroU64::new(max).expect("test message limit should be non-zero")
}
fn temp_dir(name: &str) -> PathBuf {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("clock")
.as_nanos();
let root = std::env::temp_dir().join(format!("alma-plugin-import-session-{name}-{nanos}"));
if Path::new(&root).exists() {
fs::remove_dir_all(&root).expect("stale temp dir");
}
fs::create_dir(&root).expect("temp dir");
root
}