Skip to content

Commit 389eeb8

Browse files
committed
fix(stream): serialize admission with catalog edits
1 parent b86c947 commit 389eeb8

3 files changed

Lines changed: 234 additions & 10 deletions

File tree

src/catalog_transaction.rs

Lines changed: 64 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -600,7 +600,7 @@ fn agent_semantic_deltas(
600600
}
601601

602602
fn normalize_agent(spec: &agent_spec::AgentSpec) -> Result<BTreeMap<String, SemanticAtom>> {
603-
use agent_spec::{Restart, RestartMode, TaskKind, TaskLifecycle};
603+
use agent_spec::{Restart, RestartMode, StreamLaunch, TaskKind, TaskLifecycle};
604604

605605
let host = spec
606606
.host
@@ -682,6 +682,48 @@ fn normalize_agent(spec: &agent_spec::AgentSpec) -> Result<BTreeMap<String, Sema
682682
spec.delivery.map(|delivery| delivery.as_str()),
683683
);
684684

685+
for stream in &spec.streams {
686+
let stream_base = format!("{base}/streams/{}", pointer_segment(&stream.name));
687+
match &stream.launch {
688+
None => insert_value(
689+
&mut fields,
690+
&format!("{stream_base}/launch"),
691+
SemanticType::String,
692+
"external",
693+
),
694+
Some(StreamLaunch::Command(command)) => {
695+
insert_value(
696+
&mut fields,
697+
&format!("{stream_base}/launch"),
698+
SemanticType::String,
699+
"command",
700+
);
701+
insert_value(
702+
&mut fields,
703+
&format!("{stream_base}/command"),
704+
SemanticType::String,
705+
command,
706+
);
707+
}
708+
Some(StreamLaunch::Argv(argv)) => {
709+
insert_value(
710+
&mut fields,
711+
&format!("{stream_base}/launch"),
712+
SemanticType::String,
713+
"argv",
714+
);
715+
for (index, argument) in argv.iter().enumerate() {
716+
insert_value(
717+
&mut fields,
718+
&format!("{stream_base}/argv/{index}"),
719+
SemanticType::String,
720+
argument,
721+
);
722+
}
723+
}
724+
}
725+
}
726+
685727
let restart = spec.restart_policy();
686728
let default_restart = Restart::default();
687729
insert_default_value(
@@ -3421,6 +3463,27 @@ fn test_forced_cross_device(_point: &str) -> std::io::Result<()> {
34213463
mod tests {
34223464
use super::*;
34233465

3466+
#[test]
3467+
fn external_streams_participate_in_semantic_catalog_projection() {
3468+
let root = tempfile::tempdir().unwrap();
3469+
let agent = root.path().join("agents/host/worker");
3470+
std::fs::create_dir_all(&agent).unwrap();
3471+
std::fs::write(
3472+
agent.join("agent.kdl"),
3473+
"agent \"worker\" {\n host \"host\"\n command \"worker\"\n stream \"webhook\" {}\n}\n",
3474+
)
3475+
.unwrap();
3476+
let discovered = agent_spec::discover_strict(root.path());
3477+
assert!(discovered.errors.is_empty(), "{:?}", discovered.errors);
3478+
3479+
let fields = normalize_agent(&discovered.specs[0]).unwrap();
3480+
3481+
let atom = fields
3482+
.get("/agents/host/worker/streams/webhook/launch")
3483+
.expect("external stream must affect semantic projection");
3484+
assert_eq!(atom, &present_atom(SemanticType::String, "external"));
3485+
}
3486+
34243487
#[test]
34253488
fn workspace_directory_facts_are_typed_into_the_projection_hash() {
34263489
let files = BTreeMap::new();

src/event.rs

Lines changed: 16 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -194,7 +194,7 @@ pub fn emit(
194194
}
195195
// Serialize the eligibility observation with self-authoring and desired-state changes. Once a
196196
// suspension edit owns this lock, no later emit can publish from a stale running observation.
197-
let _catalog_lock = crate::catalog_lock::CatalogLock::exclusive(root)?;
197+
let _catalog_lock = crate::catalog_lock::CatalogLock::shared(root)?;
198198
let resolved = resolve_stream(root, this_host, recipient, stream)?;
199199
let canonical_recipient = resolved.recipient;
200200
let from = format!("{canonical_recipient}/{stream}");
@@ -305,11 +305,13 @@ pub fn emit(
305305
root,
306306
&canonical_recipient,
307307
this_host,
308-
|inbox, _archive| {
308+
|inbox, archive| {
309309
for entry in record.recent.iter().filter(|entry| {
310310
key.is_none_or(|key| entry.key.as_deref() == Some(key))
311311
}) {
312-
if read_message_entry(inbox, &entry.filename)?.is_some() {
312+
if read_message_entry(inbox, &entry.filename)?.is_some()
313+
&& read_message_entry(archive, &entry.filename)?.is_none()
314+
{
313315
return Ok(Some(entry.clone()));
314316
}
315317
}
@@ -353,7 +355,17 @@ pub fn emit(
353355
);
354356
false
355357
}
356-
None => message::materialize_message_once(inbox, &filename, &rendered)?,
358+
None => match read_message_entry(inbox, &filename)? {
359+
Some(bytes) => {
360+
validate_pending_message(
361+
stream,
362+
record.pending.as_ref().expect("pending reservation exists"),
363+
&bytes,
364+
)?;
365+
false
366+
}
367+
None => message::materialize_message_once(inbox, &filename, &rendered)?,
368+
},
357369
};
358370
test_event_checkpoint(event_id, "materialized")?;
359371
if let Some(predecessor) = predecessor.as_ref() {

tests/event_e2e.rs

Lines changed: 154 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,8 @@
11
use std::fs;
22
use std::path::{Path, PathBuf};
33
use std::process::{Command, Stdio};
4-
use std::sync::{Arc, Barrier, Mutex};
4+
use std::sync::{Arc, Barrier, Mutex, mpsc};
5+
use std::time::Duration;
56

67
use sha2::{Digest as _, Sha256};
78
use st2::event::{self, EventReceiptStatus, RING_CAPACITY};
@@ -412,7 +413,155 @@ fn supersede_skips_an_archived_head_and_retires_the_latest_unread_predecessor()
412413
}
413414

414415
#[test]
415-
fn failed_predecessor_archive_leaves_the_successor_unread_and_replay_completes() {
416+
fn supersede_skips_an_archive_shadowed_inbox_candidate() {
417+
let catalog = tempfile::tempdir().unwrap();
418+
let agent = declare_agent(catalog.path(), "\"running\"", " stream \"gh-ci\" {}\n");
419+
let older = emit(catalog.path(), "pr1-queued", Some("pr-1"), false);
420+
let shadowed = emit(catalog.path(), "pr1-running", Some("pr-1"), false);
421+
let inbox = message::inbox_dir(&agent);
422+
let archive = message::archive_dir(&agent);
423+
fs::create_dir_all(&archive).unwrap();
424+
fs::copy(
425+
inbox.join(&shadowed.filename),
426+
archive.join(&shadowed.filename),
427+
)
428+
.unwrap();
429+
430+
let successor = emit(catalog.path(), "pr1-pass", Some("pr-1"), true);
431+
432+
assert_eq!(
433+
successor.superseded.as_deref(),
434+
Some(older.filename.as_str())
435+
);
436+
assert!(inbox.join(&shadowed.filename).is_file());
437+
assert!(archive.join(&shadowed.filename).is_file());
438+
assert!(archive.join(&older.filename).is_file());
439+
}
440+
441+
#[test]
442+
fn initial_supersession_authenticates_predecessor_immediately_before_archive() {
443+
let catalog = tempfile::tempdir().unwrap();
444+
let agent = declare_agent(catalog.path(), "\"running\"", " stream \"gh-ci\" {}\n");
445+
let predecessor = emit(catalog.path(), "running", Some("pr-1"), false);
446+
let inbox = message::inbox_dir(&agent);
447+
fs::write(inbox.join(&predecessor.filename), "forged predecessor").unwrap();
448+
449+
let error = event::emit(
450+
catalog.path(),
451+
"hetz",
452+
"hetz.worker",
453+
"gh-ci",
454+
"passed",
455+
Some("pr-1"),
456+
Some("passed"),
457+
"passed",
458+
true,
459+
)
460+
.unwrap_err();
461+
462+
assert!(error.to_string().contains("different bytes"), "{error:#}");
463+
assert!(inbox.join(&predecessor.filename).is_file());
464+
assert!(
465+
!message::archive_dir(&agent)
466+
.join(&predecessor.filename)
467+
.exists()
468+
);
469+
assert!(
470+
message::list_inbox(&inbox)
471+
.unwrap()
472+
.iter()
473+
.any(|message| { message.event_id.as_deref() == Some("passed") })
474+
);
475+
}
476+
477+
#[test]
478+
fn catalog_authoring_and_emit_linearize_without_deadlock() {
479+
let catalog = tempfile::tempdir().unwrap();
480+
declare_agent(catalog.path(), "\"running\"", " stream \"gh-ci\" {}\n");
481+
let admission = st2::CatalogLock::shared(catalog.path()).unwrap();
482+
let root = catalog.path().to_path_buf();
483+
let (done_tx, done_rx) = mpsc::channel();
484+
let author = std::thread::spawn(move || {
485+
let result = st2::agent_author::remove_stream(&root, "hetz.worker", "hetz", None, "gh-ci");
486+
done_tx.send(result).unwrap();
487+
});
488+
assert!(done_rx.recv_timeout(Duration::from_millis(100)).is_err());
489+
490+
let receipt = emit(catalog.path(), "before-remove", None, false);
491+
assert_eq!(receipt.status, EventReceiptStatus::Created);
492+
drop(admission);
493+
done_rx
494+
.recv_timeout(Duration::from_secs(5))
495+
.unwrap()
496+
.unwrap();
497+
author.join().unwrap();
498+
499+
let error = event::emit(
500+
catalog.path(),
501+
"hetz",
502+
"hetz.worker",
503+
"gh-ci",
504+
"after-remove",
505+
None,
506+
None,
507+
"after",
508+
false,
509+
)
510+
.unwrap_err();
511+
assert!(
512+
error.to_string().contains("does not declare stream"),
513+
"{error:#}"
514+
);
515+
}
516+
517+
#[test]
518+
fn desired_state_authoring_and_emit_linearize_without_deadlock() {
519+
let catalog = tempfile::tempdir().unwrap();
520+
declare_agent(catalog.path(), "\"running\"", " stream \"gh-ci\" {}\n");
521+
let admission = st2::CatalogLock::shared(catalog.path()).unwrap();
522+
let root = catalog.path().to_path_buf();
523+
let (done_tx, done_rx) = mpsc::channel();
524+
let author = std::thread::spawn(move || {
525+
let result = st2::agent_author::set_desired_state(
526+
&root,
527+
"hetz.worker",
528+
"hetz",
529+
None,
530+
st2::agent_author::DesiredStateValue::Suspended,
531+
Some("maintenance"),
532+
);
533+
done_tx.send(result).unwrap();
534+
});
535+
assert!(done_rx.recv_timeout(Duration::from_millis(100)).is_err());
536+
537+
assert_eq!(
538+
emit(catalog.path(), "before-suspend", None, false).status,
539+
EventReceiptStatus::Created
540+
);
541+
drop(admission);
542+
done_rx
543+
.recv_timeout(Duration::from_secs(5))
544+
.unwrap()
545+
.unwrap();
546+
author.join().unwrap();
547+
548+
let error = event::emit(
549+
catalog.path(),
550+
"hetz",
551+
"hetz.worker",
552+
"gh-ci",
553+
"after-suspend",
554+
None,
555+
None,
556+
"after",
557+
false,
558+
)
559+
.unwrap_err();
560+
assert!(error.to_string().contains("is suspended"), "{error:#}");
561+
}
562+
563+
#[test]
564+
fn invalid_archive_receipt_blocks_supersession_before_successor_publication() {
416565
let catalog = tempfile::tempdir().unwrap();
417566
let agent = declare_agent(catalog.path(), "\"running\"", " stream \"gh-ci\" {}\n");
418567
let predecessor = emit(catalog.path(), "pr1-running", Some("pr-1"), false);
@@ -438,17 +587,17 @@ fn failed_predecessor_archive_leaves_the_successor_unread_and_replay_completes()
438587
"{error:#}"
439588
);
440589
let unread = message::list_inbox(&inbox).unwrap();
441-
assert_eq!(unread.len(), 2);
590+
assert_eq!(unread.len(), 1);
442591
assert!(
443592
unread
444593
.iter()
445-
.any(|message| message.event_id.as_deref() == Some("pr1-pass"))
594+
.all(|message| message.event_id.as_deref() != Some("pr1-pass"))
446595
);
447596
assert!(inbox.join(&predecessor.filename).exists());
448597

449598
fs::remove_dir(archive.join(&predecessor.filename)).unwrap();
450599
let replay = emit(catalog.path(), "pr1-pass", Some("pr-1"), true);
451-
assert_eq!(replay.status, EventReceiptStatus::Deduplicated);
600+
assert_eq!(replay.status, EventReceiptStatus::Created);
452601
assert_eq!(
453602
replay.superseded.as_deref(),
454603
Some(predecessor.filename.as_str())

0 commit comments

Comments
 (0)