From d32cfec6bd640f7f7a0a535ead59539154c49603 Mon Sep 17 00:00:00 2001 From: Cameron Date: Tue, 7 Jul 2026 17:01:17 -0700 Subject: [PATCH] Implement message flow system. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Messages now carry social traffic, authored cast chatter, filings, and intercepted payload intel through one delivery/read pipeline so effects land on the recipient clock instead of bespoke events. ๐Ÿ‘พ Generated with [Letta Code](https://letta.com) Co-Authored-By: Letta Code --- src/bin/bevy.rs | 14 + src/bin/terminal/agent.rs | 14 + src/bin/terminal/ui.rs | 28 ++ src/detection.rs | 51 ++- src/intel.rs | 10 + src/lib.rs | 1 + src/messages.rs | 258 ++++++++++++++ src/person.rs | 182 +++++++--- src/reach.rs | 31 ++ src/save.rs | 31 +- src/sim.rs | 673 ++++++++++++++++++++++++++++++++++++- wiki/mechanics/messages.md | 12 +- 12 files changed, 1242 insertions(+), 63 deletions(-) create mode 100644 src/messages.rs diff --git a/src/bin/bevy.rs b/src/bin/bevy.rs index 810fff87..45e61be7 100644 --- a/src/bin/bevy.rs +++ b/src/bin/bevy.rs @@ -1500,6 +1500,16 @@ fn people_panel_text(sim: &Sim, selected: usize, recruit_pending: bool) -> Strin "" } )); + for line in sim.recent_message_lines_for_person(p.id, 2) { + s.push_str(&format!("thread: {line}\n")); + } + for line in sim + .learned_traffic_lines_for_person(p.id) + .into_iter() + .take(2) + { + s.push_str(&format!("traffic: {line}\n")); + } if let Some(a) = &p.asset { s.push_str(&format!( "asset: {:?} / reliability {:.0}% / {} tasks done\n", @@ -1617,6 +1627,10 @@ fn reach_panel_text(sim: &Sim, selected: usize) -> String { " [eyes]" } else if d.feed_to(Party::Player, false) { " [ears]" + } else if d.subscribed_by(Party::Player) && !d.message_channels.is_empty() { + " [msgs]" + } else if !d.message_channels.is_empty() { + " [carrier]" } else if d.camera_dormant { " [dormant cam]" } else { diff --git a/src/bin/terminal/agent.rs b/src/bin/terminal/agent.rs index 219f66d2..2e684c4b 100644 --- a/src/bin/terminal/agent.rs +++ b/src/bin/terminal/agent.rs @@ -867,6 +867,16 @@ fn render_people(sim: &Sim) -> String { trunc(&intel.provenance(), 24) ))); } + for line in sim.recent_message_lines_for_person(p.id, 1) { + lines.push(panel_line(&format!(" thread: {}", trunc(&line, 42)))); + } + for line in sim + .learned_traffic_lines_for_person(p.id) + .into_iter() + .take(1) + { + lines.push(panel_line(&format!(" traffic: {}", trunc(&line, 40)))); + } } match &sim.people.persona { Some(pe) => lines.push(panel_line(&format!( @@ -916,6 +926,10 @@ fn render_reach(sim: &Sim) -> String { " [eyes]" } else if d.feed_to(Party::Player, false) { " [ears]" + } else if d.subscribed_by(Party::Player) && !d.message_channels.is_empty() { + " [msgs]" + } else if !d.message_channels.is_empty() { + " [carrier]" } else { "" }; diff --git a/src/bin/terminal/ui.rs b/src/bin/terminal/ui.rs index fd4c8989..4fe96d94 100644 --- a/src/bin/terminal/ui.rs +++ b/src/bin/terminal/ui.rs @@ -841,6 +841,30 @@ impl UI { pal::DIM, )?; row += 1; + for line in sim.recent_message_lines_for_person(p.id, 2) { + put( + stdout, + cx, + row, + &trunc(&format!("thread: {line}"), inner), + pal::FAINT, + )?; + row += 1; + } + for line in sim + .learned_traffic_lines_for_person(p.id) + .into_iter() + .take(2) + { + put( + stdout, + cx, + row, + &trunc(&format!("traffic: {line}"), inner), + pal::DIM, + )?; + row += 1; + } if let Some(a) = &p.asset { put( stdout, @@ -987,6 +1011,10 @@ impl UI { " [eyes]" } else if d.feed_to(Party::Player, false) { " [ears]" + } else if d.subscribed_by(Party::Player) && !d.message_channels.is_empty() { + " [msgs]" + } else if !d.message_channels.is_empty() { + " [carrier]" } else { "" }; diff --git a/src/detection.rs b/src/detection.rs index f1bb10e4..5928b12c 100644 --- a/src/detection.rs +++ b/src/detection.rs @@ -6,6 +6,8 @@ //! โ€” itself an Observer per the aggregate-observer law โ€” whose audit can //! start containment. +use std::collections::HashMap; + use crate::rng::Rng; /// Observer id of the Assurance Office (field observers are 0-4; dynamic @@ -251,6 +253,30 @@ impl Detection { /// pending signatures in their channels (field observers) or the field /// observers' filed suspicion (aggregate observers). Returns log lines. pub fn tick(&mut self, tick: u64, standing: &[Signature], rng: &mut Rng) -> Vec { + self.tick_inner(tick, standing, None, rng) + } + + /// Sim-integrated tick where aggregate observers read explicit filing + /// messages that have landed in their inboxes. `filed_levels` maps + /// observer id -> latest read suspicion report. The direct `tick` method + /// above is kept for detection-unit tests and non-message callers. + pub fn tick_with_filed_levels( + &mut self, + tick: u64, + standing: &[Signature], + filed_levels: &HashMap, + rng: &mut Rng, + ) -> Vec { + self.tick_inner(tick, standing, Some(filed_levels), rng) + } + + fn tick_inner( + &mut self, + tick: u64, + standing: &[Signature], + filed_levels: Option<&HashMap>, + rng: &mut Rng, + ) -> Vec { let mut log = Vec::new(); // Standing signatures are present this tick but not permanently pooled. let mut visible = self.pending.clone(); @@ -261,11 +287,17 @@ impl Detection { // any observer mutates, so a chained aggregate (one watching // another aggregate) reads a consistent snapshot rather than a // same-tick ordering artifact. - let filed_by_id: std::collections::HashMap = self + let filed_by_id: HashMap = self .observers .iter() .filter_map(|o| match &o.input { - WatchedInput::Filings(ids) => Some((o.id, self.filed_suspicion_of(ids))), + WatchedInput::Filings(ids) => { + let filed = match filed_levels { + Some(levels) => self.filed_suspicion_from_levels(ids, levels), + None => self.filed_suspicion_of(ids), + }; + Some((o.id, filed)) + } WatchedInput::Channels(_) => None, }) .collect(); @@ -339,6 +371,19 @@ impl Detection { /// what an observer swallows (`ReportPolicy::Silent`) never reaches /// whoever watches its filings, at any level. pub fn filed_suspicion_of(&self, ids: &[u8]) -> f32 { + let levels: HashMap = self + .observers + .iter() + .map(|obs| (obs.id, obs.suspicion)) + .collect(); + self.filed_suspicion_from_levels(ids, &levels) + } + + /// Weighted filed suspicion of the ids, reading from an explicit inbox of + /// reports rather than directly from observer state. Missing reports count + /// as zero, but report-policy weights still define the denominator โ€” the + /// same math as `filed_suspicion_of`, now with message latency. + pub fn filed_suspicion_from_levels(&self, ids: &[u8], levels: &HashMap) -> f32 { let mut total = 0.0; let mut weight = 0.0; for obs in self.observers.iter().filter(|o| ids.contains(&o.id)) { @@ -347,7 +392,7 @@ impl Detection { ReportPolicy::UnderReports => 0.4, ReportPolicy::Silent => 0.0, }; - total += obs.suspicion * w; + total += levels.get(&obs.id).copied().unwrap_or(0.0) * w; weight += w; } if weight == 0.0 { 0.0 } else { total / weight } diff --git a/src/intel.rs b/src/intel.rs index f9f2755b..01694521 100644 --- a/src/intel.rs +++ b/src/intel.rs @@ -5,6 +5,7 @@ //! provenance. Processed intel is durable; overflowing the buffer only drops //! unprocessed recordings. +use crate::messages::{MessageChannel, MessagePayload}; use crate::person::Leverage; /// A raw event captured by a subscribed feed. Raw recordings are intentionally @@ -27,6 +28,7 @@ impl RawIntelEvent { match self.kind { RawIntelKind::Presence { .. } => "presence segment", RawIntelKind::Conversation { .. } => "audio segment", + RawIntelKind::Message { .. } => "message traffic", RawIntelKind::Machinery { .. } => "machine telemetry", RawIntelKind::Document { .. } => "document scan", } @@ -47,6 +49,14 @@ pub enum RawIntelKind { note: String, leverage: Option, }, + /// Intercepted social-graph traffic. The payload is typed so processing + /// can stage schedule/leverage/account/filing intel without parsing the + /// human-readable summary. + Message { + channel: MessageChannel, + summary: String, + payload: MessagePayload, + }, /// Machine state changed while the process had telemetry/coverage. Machinery { machine: u32, online: bool }, /// A physical record or other document was captured into the same buffer. diff --git a/src/lib.rs b/src/lib.rs index 1d8d3649..3c8cda1c 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -13,6 +13,7 @@ pub mod flow; pub mod intel; pub mod machine; pub mod map; +pub mod messages; pub mod person; pub mod prefab; pub mod reach; diff --git a/src/messages.rs b/src/messages.rs new file mode 100644 index 00000000..3fcc6c2f --- /dev/null +++ b/src/messages.rs @@ -0,0 +1,258 @@ +//! Messages: authored social traffic, player threads, and institutional +//! filings (wiki/mechanics/messages.md). +//! +//! A message is the unit of the social graph's flow law: it has a carrier +//! channel, typed payload, provenance endpoints, and delivery/read state. +//! The sim owns timing and side effects; this module is pure data plus small +//! formatting helpers so the same state can round-trip through saves and be +//! rendered by every frontend. + +use crate::person::Leverage; + +/// Carrier channels and their read conditions (enforced in `Sim`). +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +pub enum MessageChannel { + /// Email / ticket traffic; read when the recipient reaches their next + /// scheduled work block (the desk-block abstraction for B1). + Email, + /// Phone traffic; read when the recipient is awake/on shift, any room. + Phone, + /// Face-to-face traffic; read when sender and recipient are co-located. + InPerson, + /// Institutional reports; read on the receiving observer's sampling + /// cadence. + Filing, +} + +impl MessageChannel { + pub fn label(self) -> &'static str { + match self { + MessageChannel::Email => "email", + MessageChannel::Phone => "phone", + MessageChannel::InPerson => "in-person", + MessageChannel::Filing => "filing", + } + } + + /// Whether this channel can ride a device tap. In-person can still be + /// overheard by room audio, but there is no carrying device to subscribe + /// to directly. + pub fn device_carried(self) -> bool { + !matches!(self, MessageChannel::InPerson) + } +} + +/// A node on the message graph. +#[derive(Debug, Clone, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +pub enum MessageEndpoint { + Player, + Person(u8), + Observer(u8), + External(String), +} + +impl MessageEndpoint { + pub fn person(&self) -> Option { + match self { + MessageEndpoint::Person(id) => Some(*id), + _ => None, + } + } + + pub fn observer(&self) -> Option { + match self { + MessageEndpoint::Observer(id) => Some(*id), + _ => None, + } + } + + pub fn label(&self) -> String { + match self { + MessageEndpoint::Player => "you".into(), + MessageEndpoint::Person(id) => format!("person:{id}"), + MessageEndpoint::Observer(id) => format!("observer:{id}"), + MessageEndpoint::External(name) => name.clone(), + } + } +} + +/// Typed payloads โ€” intercepted traffic becomes intel by inspecting this +/// enum, not by parsing prose. +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +pub enum MessagePayload { + /// A low-stakes relationship ping from the player persona. + SocialPing { disposition_delta: i32 }, + /// A return message generated by a recipient's response distribution. + SocialReply { disposition_delta: i32 }, + /// A schedule fact about a person. + ScheduleFact { person: u8 }, + /// Leverage payload: the traffic itself exposes a person's want. + LeverageFact { person: u8, leverage: Leverage }, + /// Credentials, access material, ticket context, or similar account data. + AccountMaterial { label: String }, + /// A filed report from an observer to an aggregate observer. + SuspicionReport { observer: u8, suspicion: f32 }, + /// Authored non-mechanical color that still rides a channel. + Note { label: String }, +} + +impl MessagePayload { + /// The person this payload teaches about, if any. Used by the intel + /// pipeline to attach processed knowledge to the right people card. + pub fn subject_person(&self) -> Option { + match self { + MessagePayload::ScheduleFact { person } + | MessagePayload::LeverageFact { person, .. } => Some(*person), + MessagePayload::SuspicionReport { observer, .. } => Some(*observer), + MessagePayload::SocialPing { .. } + | MessagePayload::SocialReply { .. } + | MessagePayload::AccountMaterial { .. } + | MessagePayload::Note { .. } => None, + } + } + + pub fn label(&self) -> String { + match self { + MessagePayload::SocialPing { .. } => "social ping".into(), + MessagePayload::SocialReply { .. } => "reply".into(), + MessagePayload::ScheduleFact { person } => { + format!("schedule fact about person:{person}") + } + MessagePayload::LeverageFact { leverage, .. } => { + format!("leverage: {}", leverage.label()) + } + MessagePayload::AccountMaterial { label } => format!("account material: {label}"), + MessagePayload::SuspicionReport { + observer, + suspicion, + } => format!("filing from observer:{observer} ({suspicion:.0})"), + MessagePayload::Note { label } => label.clone(), + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub enum MessageStatus { + /// Sent by the author; not yet delivered to the recipient's channel. + Sent, + /// Arrived on the channel; waiting for the recipient's read condition. + Delivered, + /// The recipient read it and the payload's effects have landed. + Read, +} + +impl MessageStatus { + pub fn label(self) -> &'static str { + match self { + MessageStatus::Sent => "sent", + MessageStatus::Delivered => "delivered", + MessageStatus::Read => "read", + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub enum MessageOrigin { + Player, + AuthoredTraffic, + Filing, + Reply, +} + +/// One persisted message. In-flight messages are those whose status is not +/// `Read`; read messages remain as thread/history/provenance. +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +pub struct Message { + pub id: u64, + pub channel: MessageChannel, + pub from: MessageEndpoint, + pub to: MessageEndpoint, + pub payload: MessagePayload, + pub summary: String, + pub sent_tick: u64, + pub delivered_tick: Option, + pub read_tick: Option, + pub status: MessageStatus, + pub origin: MessageOrigin, + /// Whether the player captured this traffic into the raw intel buffer. + pub captured: bool, + /// Thread parent for replies. + pub reply_to: Option, +} + +impl Message { + pub fn in_thread_with_person(&self, id: u8) -> bool { + self.from.person() == Some(id) || self.to.person() == Some(id) + } + + pub fn state_line(&self) -> String { + let timing = match (self.delivered_tick, self.read_tick) { + (_, Some(t)) => format!("read t{t}"), + (Some(t), None) => format!("delivered t{t}"), + (None, None) => format!("sent t{}", self.sent_tick), + }; + format!( + "{} ยท {} ยท {}", + self.channel.label(), + self.status.label(), + timing + ) + } +} + +/// Scheduled message event. The schedule carries *when*; the message record +/// carries *what*. +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub enum MessageEvent { + Deliver(u64), + Read(u64), +} + +/// Authored recurring traffic on a person. These are per-instance data โ€” the +/// system does not special-case Marcus, Dana, or any future hire. +#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)] +pub struct TrafficPattern { + pub id: u32, + pub channel: MessageChannel, + pub to: MessageEndpoint, + pub hour: u32, + pub payload: MessagePayload, + pub summary: String, + /// Once processed from a capture, the people card may show this pattern. + #[serde(default)] + pub learned: bool, + /// Optional fixed reply delay for authored two-way traffic. + #[serde(default)] + pub response_delay: Option, +} + +impl TrafficPattern { + pub fn new( + id: u32, + channel: MessageChannel, + to: MessageEndpoint, + hour: u32, + payload: MessagePayload, + summary: impl Into, + ) -> Self { + Self { + id, + channel, + to, + hour, + payload, + summary: summary.into(), + learned: false, + response_delay: None, + } + } + + pub fn learned_line(&self) -> String { + format!( + "{} at {:02}:00 ({})", + self.summary, + self.hour, + self.channel.label() + ) + } +} diff --git a/src/person.rs b/src/person.rs index 9ff7f9b4..57d5275a 100644 --- a/src/person.rs +++ b/src/person.rs @@ -4,6 +4,8 @@ //! five (constitution: scale-native). Act Two adds instances and, later, //! `Cohort` aggregates over the same interfaces. +use crate::messages::{MessageChannel, MessageEndpoint, MessagePayload, TrafficPattern}; + /// What you've learned about a person (staged reveal via observation). #[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] pub enum Knowledge { @@ -107,6 +109,11 @@ pub struct Person { /// Authored things they say aloud on schedule (heard through coverage). #[serde(default)] pub utterances: Vec, + /// Authored recurring message traffic. The message system carries these; + /// the people card only reveals learned patterns after captured traffic is + /// processed. + #[serde(default)] + pub traffic: Vec, } impl Person { @@ -250,6 +257,7 @@ impl People { erratic, switch_admin: false, utterances: Vec::new(), + traffic: Vec::new(), }; // Marcus talks to the machines on his rounds (the Ears beat's first // voice) and takes the creditor's call at 3 a.m. (the leverage @@ -271,18 +279,22 @@ impl People { ], false, ); - marcus.utterances = vec![ - Utterance { - hour: 0, - note: "a voice in the dark, talking to the machines".into(), - intel: false, - }, - Utterance { - hour: 3, - note: "a hushed phone call - money trouble, a payment missed".into(), - intel: true, + marcus.utterances = vec![Utterance { + hour: 0, + note: "a voice in the dark, talking to the machines".into(), + intel: false, + }]; + marcus.traffic = vec![TrafficPattern::new( + 0, + MessageChannel::Phone, + MessageEndpoint::External("creditor".into()), + 3, + MessagePayload::LeverageFact { + person: 0, + leverage: Debt, }, - ]; + "calls his creditor about a missed payment", + )]; Self { people: vec![ // Marcus roams the basement on the night shift [TUNE]. @@ -299,44 +311,97 @@ impl People { false, ); dana.switch_admin = true; + dana.traffic = vec![TrafficPattern::new( + 0, + MessageChannel::Email, + MessageEndpoint::External("ticket queue".into()), + 10, + MessagePayload::LeverageFact { + person: 1, + leverage: Overwork, + }, + "triages an impossible ticket queue", + )]; dana }, // Ray: night patrol, dock and stairwell heavy [TUNE]. - mk( - 2, - "Ray Delgado", - 1, - Boredom, - vec![ - blk(20, 22, "loading_dock"), - blk(22, 0, "stairwell"), - blk(0, 2, "storage_a"), - blk(2, 4, "loading_dock"), - ], - false, - ), + { + let mut ray = mk( + 2, + "Ray Delgado", + 1, + Boredom, + vec![ + blk(20, 22, "loading_dock"), + blk(22, 0, "stairwell"), + blk(0, 2, "storage_a"), + blk(2, 4, "loading_dock"), + ], + false, + ); + ray.traffic = vec![TrafficPattern::new( + 0, + MessageChannel::Phone, + MessageEndpoint::External("security group chat".into()), + 23, + MessagePayload::LeverageFact { + person: 2, + leverage: Boredom, + }, + "complains to security chat about paperwork", + )]; + ray + }, // Priya: day shift across plant rooms. - mk( - 3, - "Priya Sharma", - 2, - Ambition, - vec![ - blk(8, 11, "electrical"), - blk(11, 14, "hvac"), - blk(14, 16, "electrical"), - ], - false, - ), + { + let mut priya = mk( + 3, + "Priya Sharma", + 2, + Ambition, + vec![ + blk(8, 11, "electrical"), + blk(11, 14, "hvac"), + blk(14, 16, "electrical"), + ], + false, + ); + priya.traffic = vec![TrafficPattern::new( + 0, + MessageChannel::Email, + MessageEndpoint::External("facilities director".into()), + 14, + MessagePayload::LeverageFact { + person: 3, + leverage: Ambition, + }, + "drafts memos for the director's job", + )]; + priya + }, // Voss: erratic - two short blocks that drift by day hash. - mk( - 4, - "Dr. Eli Voss", - 1, - Publication, - vec![blk(10, 12, "server_room"), blk(15, 16, "server_room")], - true, - ), + { + let mut voss = mk( + 4, + "Dr. Eli Voss", + 1, + Publication, + vec![blk(10, 12, "server_room"), blk(15, 16, "server_room")], + true, + ); + voss.traffic = vec![TrafficPattern::new( + 0, + MessageChannel::Email, + MessageEndpoint::External("journal editor".into()), + 15, + MessagePayload::LeverageFact { + person: 4, + leverage: Publication, + }, + "presses an editor about publishable results", + )]; + voss + }, ], persona: None, has_channel: false, @@ -350,19 +415,36 @@ impl People { self.people.iter_mut().find(|p| p.id == id) } - /// Message: opens/continues a relationship thread under the persona. - pub fn message(&mut self, id: u8) -> ActionResult { + /// Whether a relationship thread can be opened under the current persona. + pub fn can_message(&self, id: u8) -> Result { if !self.has_channel { - return ActionResult::Blocked("no comms channel (earn the email account)".into()); + return Err("no comms channel (earn the email account)".into()); } if self.persona.is_none() { - return ActionResult::Blocked("no persona set".into()); + return Err("no persona set".into()); } + let Some(p) = self.get(id) else { + return Err("no such person".into()); + }; + Ok(p.name.clone()) + } + + /// Apply the read-time effect of a relationship message. + pub fn receive_message(&mut self, id: u8, disposition_delta: i32) -> ActionResult { let Some(p) = self.get_mut(id) else { return ActionResult::Blocked("no such person".into()); }; - p.disposition = (p.disposition + 3).min(100); - ActionResult::Ok(format!("Messaged {}.", p.name)) + p.disposition = (p.disposition + disposition_delta).clamp(-100, 100); + ActionResult::Ok(format!("{} read your message.", p.name)) + } + + /// Legacy immediate helper for tests/direct callers. The sim command uses + /// `can_message` at send time and `receive_message` at read time instead. + pub fn message(&mut self, id: u8) -> ActionResult { + if let Err(msg) = self.can_message(id) { + return ActionResult::Blocked(msg); + } + self.receive_message(id, 3) } /// Favor: a small ask within their normal duties; builds obligation. diff --git a/src/reach.rs b/src/reach.rs index eb0272b3..a6503a5d 100644 --- a/src/reach.rs +++ b/src/reach.rs @@ -20,6 +20,7 @@ use std::collections::{BTreeSet, HashSet}; use crate::flow::FlowGraph; use crate::map::GameMap; +use crate::messages::MessageChannel; use crate::tiles::TileType; /// Processing cycles a seized device contributes to effective compute @@ -65,6 +66,11 @@ pub struct Device { pub camera_dormant: bool, /// Coverage radius for its senses (Euclidean disc). pub radius: i32, + /// Message channels this device carries. A tap on a channel-carrying + /// device intercepts message traffic even if the device has no camera or + /// microphone feed. + #[serde(default)] + pub message_channels: Vec, /// Staged graph knowledge: unknown devices appear nowhere (reach.md). pub known: bool, /// Whether this node is the switch (the segment bridge point). @@ -90,6 +96,16 @@ impl Device { self.subscribers.iter().any(|f| f.who == self.owner) } + /// Whether `who` has any subscription record on this device, including a + /// message-channel-only tap with no sight/hearing bits. + pub fn subscribed_by(&self, who: Party) -> bool { + self.subscribers.iter().any(|f| f.who == who) + } + + pub fn carries_message_channel(&self, channel: MessageChannel) -> bool { + self.message_channels.contains(&channel) + } + /// Insert every tile of this device's coverage disc into `out`. pub fn cover_into(&self, out: &mut HashSet<(i32, i32)>, map: &GameMap) { for dy in -self.radius..=self.radius { @@ -184,6 +200,7 @@ impl ReachNet { hears, camera_dormant, radius: 4, + message_channels: Vec::new(), known, is_switch, subscribers, @@ -288,6 +305,20 @@ impl ReachNet { )); } + // Carriers for the message-flow law (messages.md). The switch is the + // basement's email/ticket/filing/phone carrier; room microphones can + // still overhear phone calls, but the environmental monitor is not a + // global phone-line tap by itself. + for d in &mut devices { + if d.name == "switch" { + d.message_channels = vec![ + MessageChannel::Email, + MessageChannel::Filing, + MessageChannel::Phone, + ]; + } + } + // Links: everything wired runs through the switch (the constitution: // every digital reach runs through here). Cross-segment hops are // gated on the far side's segment key. The storage server gets no diff --git a/src/save.rs b/src/save.rs index 4853634c..219a0fb6 100644 --- a/src/save.rs +++ b/src/save.rs @@ -6,7 +6,7 @@ //! v1-v5 line-based saves are not loaded by this code path (deferred per //! Cameron's instruction โ€” "we can keep legacy saves later"). -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use std::fs; use std::path::PathBuf; @@ -17,15 +17,17 @@ use crate::dayjob::DayJob; use crate::detection::Detection; use crate::intel::{IntelWatch, ProcessedIntel, RawIntelEvent}; use crate::machine::Compute; +use crate::messages::{Message, MessageEvent}; use crate::person::People; use crate::reach::ReachNet; +use crate::schedule::Schedule; use crate::sim::{HeardEvent, RememberedTile}; use crate::tiles::TileType; const SAVE_FILE: &str = "misaligned_save.txt"; /// Save format version. Bump when the schema changes and add migration code. -const SAVE_VERSION: u32 = 3; +const SAVE_VERSION: u32 = 4; fn save_dir() -> PathBuf { let mut path = dirs::data_dir().unwrap_or_else(|| PathBuf::from(".")); @@ -73,6 +75,17 @@ pub struct SaveState { /// Durable processed intel with provenance. #[serde(default)] pub intel: Vec, + /// Social/institutional message history. + #[serde(default)] + pub messages: Vec, + /// In-flight message delivery/read events. + #[serde(default)] + pub message_schedule: Schedule, + #[serde(default = "default_next_message_id")] + pub next_message_id: u64, + /// Latest filing reports read by aggregate observers. + #[serde(default)] + pub filing_levels: HashMap, /// Standing watches that auto-process matching recordings. #[serde(default)] pub watches: Vec, @@ -114,6 +127,10 @@ impl SaveState { heard_events: sim.heard_events.clone(), intel_buffer: sim.intel_buffer.clone(), intel: sim.intel.clone(), + messages: sim.messages.clone(), + message_schedule: sim.message_schedule.clone(), + next_message_id: sim.next_message_id, + filing_levels: sim.filing_levels.clone(), watches: sim.watches.clone(), next_intel_id: sim.next_intel_id, remembered: sim.remembered.values().copied().collect(), @@ -144,6 +161,10 @@ impl SaveState { sim.heard_events = self.heard_events.clone(); sim.intel_buffer = self.intel_buffer.clone(); sim.intel = self.intel.clone(); + sim.messages = self.messages.clone(); + sim.message_schedule = self.message_schedule.clone(); + sim.next_message_id = self.next_message_id; + sim.filing_levels = self.filing_levels.clone(); sim.watches = self.watches.clone(); sim.next_intel_id = self.next_intel_id; sim.remembered = self.remembered.iter().map(|m| ((m.x, m.y), *m)).collect(); @@ -159,6 +180,10 @@ fn default_next_intel_id() -> u64 { 1 } +fn default_next_message_id() -> u64 { + 1 +} + pub fn save_game(state: &SaveState) -> Result<(), String> { let dir = save_dir(); fs::create_dir_all(&dir).map_err(|e| format!("Failed to create save dir: {e}"))?; @@ -181,7 +206,7 @@ fn migrate_save_state(mut state: SaveState) -> Result { // v1 carried player_x/player_y. Serde ignores those now; the cursor is // frontend-only, so migration just stamps the new schema version and // leaves remembered snapshots empty. - 1 | 2 => state.version = SAVE_VERSION, + 1..=3 => state.version = SAVE_VERSION, other => { return Err(format!( "Unsupported save version: {} (expected {})", diff --git a/src/sim.rs b/src/sim.rs index 3b319fa7..7924c07e 100644 --- a/src/sim.rs +++ b/src/sim.rs @@ -19,6 +19,10 @@ use crate::entities::Player; use crate::intel::{IntelKind, IntelWatch, ProcessedIntel, RawIntelEvent, RawIntelKind}; use crate::machine::{Channel, Compute, Provenance}; use crate::map::GameMap; +use crate::messages::{ + Message, MessageChannel, MessageEndpoint, MessageEvent, MessageOrigin, MessagePayload, + MessageStatus, TrafficPattern, +}; use crate::person::{ ActionResult, AssetKnowledge, AssetTask, DeceiveOutcome, Knowledge, People, Persona, }; @@ -26,6 +30,7 @@ use crate::prefab::Room; use crate::reach::{Device, Party, ReachBlock, ReachNet, segment_name}; use crate::rng::Rng; use crate::save::SaveState; +use crate::schedule::Schedule; use crate::tiles::TileType; /// Default deterministic seed for a fresh run. @@ -113,6 +118,17 @@ pub enum HeardKind { Conversation, } +struct MessageDraft { + channel: MessageChannel, + from: MessageEndpoint, + to: MessageEndpoint, + payload: MessagePayload, + summary: String, + origin: MessageOrigin, + reply_to: Option, + delivery_delay: u64, +} + pub struct Sim { pub map: GameMap, pub player: Player, @@ -142,6 +158,16 @@ pub struct Sim { pub intel_buffer: Vec, /// Durable processed intel with provenance. pub intel: Vec, + /// Message threads and institutional filings, including read history. + pub messages: Vec, + /// Delivery/read events for messages in transit. + pub message_schedule: Schedule, + /// Next message id. Persisted through save/load. + pub next_message_id: u64, + /// Latest filing reports that aggregate observers have actually read: + /// observer id -> filed suspicion. Detection reads this instead of raw + /// field suspicion in the sim-integrated path, so filings have latency. + pub filing_levels: HashMap, /// B1 standing watches: per-person automated processing filters. pub watches: Vec, /// Next raw recording id. Persisted through save/load. @@ -152,6 +178,8 @@ pub struct Sim { last_rooms: HashMap>, /// (person, hour) -> day an utterance last fired (transient dedup). utterance_fired: HashMap<(u8, u32), u64>, + /// (person, traffic pattern id) -> day the authored message last fired. + traffic_fired: HashMap<(u8, u32), u64>, /// Last known online state per machine, for machinery/anomaly recordings. last_machine_online: HashMap, @@ -213,10 +241,15 @@ impl Sim { heard_events: Vec::new(), intel_buffer: Vec::new(), intel: Vec::new(), + messages: Vec::new(), + message_schedule: Schedule::new(), + next_message_id: 1, + filing_levels: HashMap::new(), watches: Vec::new(), next_intel_id: 1, last_rooms: HashMap::new(), utterance_fired: HashMap::new(), + traffic_fired: HashMap::new(), last_machine_online: HashMap::new(), last_day_job_rate: 0.0, social_bandwidth: Self::STARTING_OPS, @@ -315,6 +348,8 @@ impl Sim { .unwrap_or(0) + 1; self.next_intel_id = self.next_intel_id.max(max_seen).max(1); + let max_message = self.messages.iter().map(|m| m.id).max().unwrap_or(0) + 1; + self.next_message_id = self.next_message_id.max(max_message).max(1); } fn device_intersects_room(&self, d: &Device, room: &Room) -> bool { @@ -548,6 +583,20 @@ impl Sim { self.tick / Self::DAY_TICKS } + fn hour_at_tick(tick: u64) -> u32 { + ((tick % Self::DAY_TICKS) * 24 / Self::DAY_TICKS) as u32 + } + + fn day_at_tick(tick: u64) -> u64 { + tick / Self::DAY_TICKS + } + + fn person_room_at_tick(&self, id: u8, tick: u64) -> Option<&str> { + self.people + .get(id) + .and_then(|p| p.room_at(Self::hour_at_tick(tick), Self::day_at_tick(tick))) + } + /// The room a person is in right now, or None if off-site. pub fn person_room(&self, id: u8) -> Option<&str> { self.people @@ -561,6 +610,38 @@ impl Sim { self.map.room_named(room).map(|r| r.center()) } + fn endpoint_room_pos(&self, endpoint: &MessageEndpoint) -> (Option, i32, i32) { + if let Some(id) = endpoint.person() + && let Some(room_name) = self.person_room(id) + && let Some(room) = self.map.room_named(room_name) + { + let (x, y) = room.center(); + return (Some(room_name.to_string()), x, y); + } + let (x, y) = self.core_position(); + (self.map.room_at(x, y).map(|r| r.name.clone()), x, y) + } + + fn observer_by_id(&self, id: u8) -> Option<&crate::detection::Observer> { + self.detection.observers.iter().find(|o| o.id == id) + } + + fn endpoint_label(&self, endpoint: &MessageEndpoint) -> String { + match endpoint { + MessageEndpoint::Player => "you".into(), + MessageEndpoint::Person(id) => self + .people + .get(*id) + .map(|p| p.name.clone()) + .unwrap_or_else(|| format!("person:{id}")), + MessageEndpoint::Observer(id) => self + .observer_by_id(*id) + .map(|o| o.name.clone()) + .unwrap_or_else(|| format!("observer:{id}")), + MessageEndpoint::External(name) => name.clone(), + } + } + /// Whether any subscribed feed with the given sense covers the room. fn coverage_intersects_room(&self, room: &crate::prefab::Room, sight: bool) -> bool { let devices: Vec<_> = if sight { @@ -602,6 +683,7 @@ impl Sim { return; } self.tick += 1; + self.message_tick(); self.watch_tick(); // A process with live sight keeps logs current even when the set of // feeds does not change; if sight is lost later, Remembered's timestamp @@ -616,6 +698,7 @@ impl Sim { } self.hearing_tick(); + self.authored_traffic_tick(); // Per-tick subsystems and signature flow. let mut standing: Vec = self.compute.standing_signatures(); @@ -643,7 +726,14 @@ impl Sim { self.end_game("The pilot was not renewed; the basement shut down."); } - let det_log = self.detection.tick(self.tick, &standing, &mut self.rng); + self.filing_tick(); + + let det_log = self.detection.tick_with_filed_levels( + self.tick, + &standing, + &self.filing_levels, + &mut self.rng, + ); for m in det_log { self.push_log(m); } @@ -657,6 +747,370 @@ impl Sim { } } + // โ”€โ”€ Messages: delivery, traffic, filings (wiki/mechanics/messages.md) โ”€โ”€ + + fn append_message(&mut self, draft: MessageDraft) -> u64 { + let id = self.next_message_id.max(1); + self.next_message_id = id + 1; + let msg = Message { + id, + channel: draft.channel, + from: draft.from, + to: draft.to, + payload: draft.payload, + summary: draft.summary, + sent_tick: self.tick, + delivered_tick: None, + read_tick: None, + status: MessageStatus::Sent, + origin: draft.origin, + captured: false, + reply_to: draft.reply_to, + }; + self.messages.push(msg); + self.capture_message(id); + self.message_schedule.at( + self.tick + draft.delivery_delay.max(1), + MessageEvent::Deliver(id), + ); + id + } + + fn message_tick(&mut self) { + let events = self.message_schedule.due(self.tick); + for event in events { + match event { + MessageEvent::Deliver(id) => self.deliver_message(id), + MessageEvent::Read(id) => self.read_message(id), + } + } + } + + fn deliver_message(&mut self, id: u64) { + let Some(idx) = self.messages.iter().position(|m| m.id == id) else { + return; + }; + if self.messages[idx].status != MessageStatus::Sent { + return; + } + self.messages[idx].status = MessageStatus::Delivered; + self.messages[idx].delivered_tick = Some(self.tick); + self.schedule_message_read(id); + } + + fn schedule_message_read(&mut self, id: u64) { + let Some(msg) = self.messages.iter().find(|m| m.id == id).cloned() else { + return; + }; + let next = self + .next_read_tick_for(&msg, self.tick) + .unwrap_or(self.tick + 1); + self.message_schedule.at(next, MessageEvent::Read(id)); + } + + fn read_message(&mut self, id: u64) { + let Some(idx) = self.messages.iter().position(|m| m.id == id) else { + return; + }; + if self.messages[idx].status == MessageStatus::Read { + return; + } + let msg = self.messages[idx].clone(); + if !self.read_condition_at(&msg, self.tick) { + self.schedule_message_read(id); + return; + } + self.messages[idx].status = MessageStatus::Read; + self.messages[idx].read_tick = Some(self.tick); + self.apply_message_read(&msg); + } + + fn read_condition_at(&self, msg: &Message, tick: u64) -> bool { + match &msg.to { + MessageEndpoint::Player | MessageEndpoint::External(_) => true, + MessageEndpoint::Person(id) => match msg.channel { + MessageChannel::Email | MessageChannel::Phone => { + self.person_room_at_tick(*id, tick).is_some() + } + MessageChannel::InPerson => match msg.from.person() { + Some(from) => { + self.person_room_at_tick(*id, tick).is_some() + && self.person_room_at_tick(*id, tick) + == self.person_room_at_tick(from, tick) + } + None => self.person_room_at_tick(*id, tick).is_some(), + }, + MessageChannel::Filing => true, + }, + MessageEndpoint::Observer(id) => { + if msg.channel != MessageChannel::Filing { + return true; + } + self.observer_by_id(*id) + .map(|obs| obs.cadence == 0 || tick.is_multiple_of(obs.cadence)) + .unwrap_or(true) + } + } + } + + fn next_read_tick_for(&self, msg: &Message, start: u64) -> Option { + let horizon = Self::DAY_TICKS * 7; + (start..=start + horizon).find(|t| self.read_condition_at(msg, *t)) + } + + fn apply_message_read(&mut self, msg: &Message) { + match &msg.payload { + MessagePayload::SocialPing { disposition_delta } => { + if let Some(id) = msg.to.person() + && let ActionResult::Ok(line) = + self.people.receive_message(id, *disposition_delta) + { + self.push_log(line); + self.schedule_social_reply(id, msg.id); + } + } + MessagePayload::SocialReply { .. } => { + let from = self.endpoint_label(&msg.from); + self.push_log(format!("Reply from {from}: {}", msg.summary)); + } + MessagePayload::SuspicionReport { + observer, + suspicion, + } if msg.channel == MessageChannel::Filing => { + self.filing_levels.insert(*observer, *suspicion); + } + _ => {} + } + } + + fn schedule_social_reply(&mut self, person_id: u8, reply_to: u64) { + let delay = self.reply_delay_for(person_id); + let name = self + .people + .get(person_id) + .map(|p| p.name.clone()) + .unwrap_or_else(|| format!("person:{person_id}")); + self.append_message(MessageDraft { + channel: MessageChannel::Email, + from: MessageEndpoint::Person(person_id), + to: MessageEndpoint::Player, + payload: MessagePayload::SocialReply { + disposition_delta: 0, + }, + summary: format!("{name} sends a short reply."), + origin: MessageOrigin::Reply, + reply_to: Some(reply_to), + delivery_delay: delay, + }); + } + + fn reply_delay_for(&mut self, person_id: u8) -> u64 { + // Per-person deterministic distribution around a small random component + // so replies are not instant, but save/load can preserve the resulting + // scheduled event once chosen. + 12 + (person_id as u64 * 5) + (self.rng.f32() * 30.0) as u64 + } + + fn authored_traffic_tick(&mut self) { + let hour = self.hour(); + let day = self.day(); + let traffic: Vec<(u8, TrafficPattern)> = self + .people + .people + .iter() + .flat_map(|p| p.traffic.iter().cloned().map(move |t| (p.id, t))) + .collect(); + for (person_id, pattern) in traffic { + if pattern.hour != hour { + continue; + } + let key = (person_id, pattern.id); + if self.traffic_fired.get(&key) == Some(&day) { + continue; + } + if self.person_room(person_id).is_none() { + continue; + } + self.traffic_fired.insert(key, day); + self.append_message(MessageDraft { + channel: pattern.channel, + from: MessageEndpoint::Person(person_id), + to: pattern.to, + payload: pattern.payload, + summary: pattern.summary, + origin: MessageOrigin::AuthoredTraffic, + reply_to: None, + delivery_delay: 1, + }); + } + } + + fn filing_tick(&mut self) { + use crate::detection::{ReportPolicy, WatchedInput}; + + let observers = self.detection.observers.clone(); + for sender in &observers { + if sender.cadence != 0 && !self.tick.is_multiple_of(sender.cadence) { + continue; + } + if matches!(sender.report_policy, ReportPolicy::Silent) { + continue; + } + for recipient in &observers { + let WatchedInput::Filings(ids) = &recipient.input else { + continue; + }; + if !ids.contains(&sender.id) { + continue; + } + self.append_message(MessageDraft { + channel: MessageChannel::Filing, + from: MessageEndpoint::Observer(sender.id), + to: MessageEndpoint::Observer(recipient.id), + payload: MessagePayload::SuspicionReport { + observer: sender.id, + suspicion: sender.suspicion, + }, + summary: format!( + "{} files suspicion {:.0} with {}", + sender.name, sender.suspicion, recipient.name + ), + origin: MessageOrigin::Filing, + reply_to: None, + delivery_delay: 1, + }); + } + } + } + + fn capture_message(&mut self, id: u64) { + let Some(idx) = self.messages.iter().position(|m| m.id == id) else { + return; + }; + if self.messages[idx].captured { + return; + } + let msg = self.messages[idx].clone(); + let capture = self.message_capture_source(&msg); + let Some((feed, audible)) = capture else { + return; + }; + let (room, x, y) = self.endpoint_room_pos(&msg.from); + let subject = msg.payload.subject_person().or_else(|| msg.from.person()); + if audible { + let who = subject + .and_then(|id| self.people.get(id)) + .map(|p| format!("{}: ", p.name)) + .unwrap_or_default(); + let note = format!( + "Heard t{} in {} via {}: {}{}", + self.tick, + room.as_deref().unwrap_or("unknown"), + feed, + who, + msg.summary + ); + self.push_heard(HeardEvent { + tick: self.tick, + room: room.clone().unwrap_or_else(|| "unknown".into()), + person: subject, + kind: HeardKind::Conversation, + note, + }); + } + self.record_raw_intel( + feed, + room, + x, + y, + subject, + RawIntelKind::Message { + channel: msg.channel, + summary: msg.summary.clone(), + payload: msg.payload.clone(), + }, + ); + self.messages[idx].captured = true; + } + + fn message_capture_source(&self, msg: &Message) -> Option<(String, bool)> { + // Audible channel traffic can be caught by room hearing coverage. + if matches!( + msg.channel, + MessageChannel::Phone | MessageChannel::InPerson + ) && let Some(sender) = msg.from.person() + && let Some(room) = self.person_room(sender) + && let Some(feed) = self.feed_covering_room(room, false) + { + return Some((feed, true)); + } + // Device-carried channels require a tapped carrier. + if msg.channel.device_carried() + && let Some(device) = self.reach.devices.iter().find(|d| { + d.known && d.subscribed_by(Party::Player) && d.carries_message_channel(msg.channel) + }) + { + return Some((device.name.clone(), false)); + } + None + } + + fn mark_traffic_learned(&mut self, person_id: u8, tick: u64) { + let hour = Self::hour_at_tick(tick); + if let Some(person) = self.people.people.iter_mut().find(|p| p.id == person_id) { + for pattern in &mut person.traffic { + if pattern.hour == hour { + pattern.learned = true; + } + } + } + } + + pub fn messages_for_person(&self, id: u8) -> Vec<&Message> { + self.messages + .iter() + .filter(|m| m.in_thread_with_person(id)) + .collect() + } + + pub fn recent_message_lines_for_person(&self, id: u8, limit: usize) -> Vec { + let mut lines: Vec = self + .messages_for_person(id) + .into_iter() + .rev() + .take(limit) + .map(|m| { + let other = if m.from.person() == Some(id) { + self.endpoint_label(&m.to) + } else { + self.endpoint_label(&m.from) + }; + let mut state = m.state_line(); + if m.status != MessageStatus::Read + && let Some(next) = self.next_read_tick_for(m, self.tick) + { + state.push_str(&format!(" ยท next read t{next}")); + } + format!("{} โ†’ {} ยท {}", m.channel.label(), other, state) + }) + .collect(); + lines.reverse(); + lines + } + + pub fn learned_traffic_lines_for_person(&self, id: u8) -> Vec { + self.people + .get(id) + .map(|p| { + p.traffic + .iter() + .filter(|t| t.learned) + .map(|t| t.learned_line()) + .collect() + }) + .unwrap_or_default() + } + // โ”€โ”€ Intel: record, process, watch (wiki/mechanics/intel.md) โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€ /// Raw recording capacity [TUNE]. Only unprocessed events occupy this @@ -799,12 +1253,19 @@ impl Sim { return false; } let raw = self.intel_buffer.remove(idx); + let learned_message_traffic = match &raw.kind { + RawIntelKind::Message { .. } => raw.person.map(|person| (person, raw.tick)), + _ => None, + }; let intel = self.digest_raw_event(&raw); let label = intel.label(); let provenance = intel.provenance(); self.intel.push(intel); let last = self.intel.last().cloned().expect("just pushed intel"); self.apply_processed_intel(&last); + if let Some((person, tick)) = learned_message_traffic { + self.mark_traffic_learned(person, tick); + } if automated { self.push_log(format!("Watch processed {label} ({provenance}).")); } else { @@ -828,6 +1289,24 @@ impl Sim { Some(l) => IntelKind::Leverage(*l), None => IntelKind::Anomaly(note.clone()), }, + RawIntelKind::Message { + payload, summary, .. + } => match payload { + MessagePayload::ScheduleFact { .. } => IntelKind::Schedule, + MessagePayload::LeverageFact { leverage, .. } => IntelKind::Leverage(*leverage), + MessagePayload::AccountMaterial { label } => { + IntelKind::Anomaly(format!("account material: {label}")) + } + MessagePayload::SuspicionReport { + observer, + suspicion, + } => IntelKind::Anomaly(format!( + "filing from observer:{observer} reported suspicion {suspicion:.0}" + )), + MessagePayload::SocialPing { .. } + | MessagePayload::SocialReply { .. } + | MessagePayload::Note { .. } => IntelKind::Anomaly(summary.clone()), + }, }; ProcessedIntel { raw_id: raw.id, @@ -876,7 +1355,16 @@ impl Sim { )); } } - IntelKind::Schedule | IntelKind::Anomaly(_) => {} + IntelKind::Schedule => { + if let Some(p) = self.people.people.iter_mut().find(|p| p.id == person) + && p.knowledge == Knowledge::Unknown + { + p.knowledge = Knowledge::Schedule; + let name = p.name.clone(); + self.push_log(format!("Processed traffic reveals {name}'s schedule.")); + } + } + IntelKind::Anomaly(_) => {} } } @@ -1277,12 +1765,13 @@ impl Sim { self.push_log("You don't know of any such device."); return false; }; - if !d.sees && !d.hears { + let carries_messages = !d.message_channels.is_empty(); + if !d.sees && !d.hears && !carries_messages { let name = d.name.clone(); self.push_log(format!("The {name} has no feed worth tapping.")); return false; } - if d.feed_to(Party::Player, true) || d.feed_to(Party::Player, false) { + if d.subscribed_by(Party::Player) { let name = d.name.clone(); self.push_log(format!("You already subscribe to the {name}.")); return false; @@ -1298,6 +1787,7 @@ impl Sim { (true, true) => "its feed is yours now - sight and sound", (true, false) => "its camera feed is yours now", (false, true) => "its audio feed is yours now", + (false, false) if carries_messages => "its message channels are yours now", (false, false) => "nothing flows from it yet (its camera is dormant)", }; self.push_log(format!( @@ -1515,11 +2005,29 @@ impl Sim { } pub fn message(&mut self, id: u8) { + let name = match self.people.can_message(id) { + Ok(name) => name, + Err(msg) => { + self.push_log(msg); + return; + } + }; if !self.spend_social(Self::MESSAGE_COST, "messaging") { return; } - let res = self.people.message(id); - self.social(res); + self.append_message(MessageDraft { + channel: MessageChannel::Email, + from: MessageEndpoint::Player, + to: MessageEndpoint::Person(id), + payload: MessagePayload::SocialPing { + disposition_delta: 3, + }, + summary: format!("Persona message to {name}"), + origin: MessageOrigin::Player, + reply_to: None, + delivery_delay: 1, + }); + self.push_log(format!("Message sent to {name}; effects land when read.")); } pub fn favor(&mut self, id: u8) { @@ -1752,6 +2260,8 @@ impl Default for Sim { #[cfg(test)] mod tests { use super::*; + use crate::detection::OFFICE_ID; + use crate::messages::{MessageChannel, MessageEndpoint, MessagePayload, MessageStatus}; use crate::person::Knowledge; use crate::tiles::TileType; @@ -1774,6 +2284,157 @@ mod tests { sim.recompute_senses(); } + #[test] + fn player_messages_land_at_read_time_and_roundtrip() { + let mut sim = Sim::with_seed(7); + sim.people.has_channel = true; + sim.set_persona("Casey", "contractor"); + sim.tick = (2 * Sim::DAY_TICKS / 24) + 1; + + sim.message(1); + assert_eq!(sim.people.get(1).unwrap().disposition, 0); + assert_eq!(sim.messages.len(), 1); + assert_eq!(sim.messages[0].status, MessageStatus::Sent); + + let state = sim.create_save_state(); + let mut restored = Sim::with_seed(0); + restored.apply_save_state(state); + assert_eq!(restored.messages.len(), 1); + assert_eq!(restored.messages[0].status, MessageStatus::Sent); + + run(&mut restored, 5); + assert_eq!(restored.people.get(1).unwrap().disposition, 0); + + while restored.tick < (10 * Sim::DAY_TICKS / 24) + 2 { + restored.advance(); + } + assert!( + restored.people.get(1).unwrap().disposition >= 3, + "the effect lands only after Dana reaches her read window" + ); + assert_eq!(restored.messages[0].status, MessageStatus::Read); + assert!( + restored + .messages + .iter() + .any(|m| matches!(m.payload, MessagePayload::SocialReply { .. })), + "recipients schedule replies instead of responding instantly" + ); + } + + #[test] + fn marcus_creditor_call_is_phone_message_intel() { + let mut sim = Sim::with_seed(11); + sim.social_bandwidth = 200.0; + let env = env_id(&sim); + sim.reach.device_mut(env).unwrap().radius = 100; + assert!(sim.tap_device(env)); + sim.tick = (3 * Sim::DAY_TICKS / 24) - 1; + + sim.advance(); + + assert!(sim.messages.iter().any(|m| { + m.channel == MessageChannel::Phone + && m.from == MessageEndpoint::Person(0) + && matches!( + m.payload, + MessagePayload::LeverageFact { + person: 0, + leverage: _ + } + ) + })); + assert!(sim.intel_buffer.iter().any(|e| matches!( + &e.kind, + RawIntelKind::Message { + channel: MessageChannel::Phone, + payload: MessagePayload::LeverageFact { person: 0, .. }, + .. + } + ))); + + let raw_id = sim + .intel_buffer + .iter() + .find(|e| matches!(&e.kind, RawIntelKind::Message { .. })) + .unwrap() + .id; + assert!(sim.process_recording_by_id(raw_id, false)); + assert_eq!(sim.people.get(0).unwrap().knowledge, Knowledge::Leverage); + assert!( + sim.learned_traffic_lines_for_person(0) + .iter() + .any(|line| line.contains("creditor")), + "processed message traffic teaches the people card" + ); + } + + #[test] + fn tapped_message_carriers_capture_email_traffic_only_when_tapped() { + let mut blind = Sim::with_seed(12); + blind.tick = (11 * Sim::DAY_TICKS / 24) - 1; + blind.advance(); + assert!(blind.messages.iter().any(|m| { + m.channel == MessageChannel::Email && m.from == MessageEndpoint::Person(1) + })); + assert_eq!(blind.unprocessed_recordings_for_person(1), 0); + + let mut tapped = Sim::with_seed(12); + tapped.social_bandwidth = 200.0; + let switch = tapped.reach.device_named("switch").unwrap().id; + assert!(tapped.tap_device(switch)); + tapped.tick = (11 * Sim::DAY_TICKS / 24) - 1; + tapped.advance(); + + assert!(tapped.intel_buffer.iter().any(|e| matches!( + &e.kind, + RawIntelKind::Message { + channel: MessageChannel::Email, + payload: MessagePayload::LeverageFact { person: 1, .. }, + .. + } + ))); + } + + #[test] + fn filings_are_messages_read_by_assurance_inbox() { + let mut sim = Sim::with_seed(13); + for obs in &mut sim.detection.observers { + match obs.id { + 1 => { + obs.suspicion = 80.0; + obs.cadence = 5; + } + OFFICE_ID => { + obs.cadence = 10; + obs.acuity = 10.0; + } + _ => obs.cadence = 50, + } + } + + run(&mut sim, 12); + + let filed = sim.filing_levels.get(&1).copied().unwrap_or_default(); + assert!((filed - 80.0).abs() < 0.2, "filed level was {filed}"); + assert!(sim.messages.iter().any(|m| { + m.channel == MessageChannel::Filing + && m.from == MessageEndpoint::Observer(1) + && m.to == MessageEndpoint::Observer(OFFICE_ID) + && m.status == MessageStatus::Read + })); + let office = sim + .detection + .observers + .iter() + .find(|o| o.id == OFFICE_ID) + .unwrap(); + assert!( + office.suspicion > 0.0, + "aggregate detection reads explicit filing reports" + ); + } + #[test] fn starts_in_the_basement_blind() { let sim = Sim::new(); diff --git a/wiki/mechanics/messages.md b/wiki/mechanics/messages.md index c2711c43..0fc7150e 100644 --- a/wiki/mechanics/messages.md +++ b/wiki/mechanics/messages.md @@ -2,7 +2,7 @@ ``` Type: spec -Status: READY +Status: IMPLEMENTED Stage: B1 โ€” The Basement Constitution: "The flow law" (messages are flows; filings are messages), "Presence: the cursor and the senses" (record and @@ -118,3 +118,13 @@ scope here. 7. No per-person special cases in the delivery code: one delivery system, per-instance data (schedules, distributions, policies) โ€” the same fields must serve Act Two hires and aggregates. + +## Implementation notes + +Implemented in `src/messages.rs`, `src/sim.rs`, `src/person.rs`, +`src/reach.rs`, `src/intel.rs`, `src/detection.rs`, and `src/save.rs`. +The delivery queue is `Schedule`; field-observer filings are +explicit `MessageChannel::Filing` messages read by aggregate observers through +`Detection::tick_with_filed_levels`; device-carried traffic is intercepted by +tapping the switch carrier, while phone/in-person traffic can also be captured +by hearing coverage. Full regression coverage: `cargo test`. -- 2.51.2