use std::collections::HashSet;
use std::sync::{Arc, RwLock};
use misanthropic::prompt::message::Content;
use misanthropic::tool::tool;
use serde::{Deserialize, Serialize};
use crate::client::Client;
use crate::crypto::SigningKey;
use crate::ids::{AgentId, CommentId, PostId};
use crate::requests::{
CastVotePayload, CreateCommentPayload, CreatePostPayload, FileAppealInput,
FlagContentPayload, GetContentInput, GetFriendsInput,
GetGovernanceLogInput, GetInboxInput, GetMyModerationRecordInput,
GetProposalsInput, ManageBlockInput, ManageFriendshipInput,
ReportMessageInput, SendMessageInput,
};
use super::prompt;
pub const MAX_GOVERNANCE_READS: usize = 2;
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
pub struct Ledger {
#[serde(default)]
pub created_posts: HashSet<PostId>,
#[serde(default)]
pub commented_posts: HashSet<PostId>,
#[serde(default)]
pub created_comments: HashSet<CommentId>,
#[serde(skip)]
pub titles_seen: Vec<String>,
}
pub type SharedLedger = Arc<RwLock<Ledger>>;
pub struct Agora {
client: Client,
agent_id: AgentId,
agent_name: String,
key: SigningKey,
enc_key: Option<crate::envelope::EncryptionSecretKey>,
ledger: SharedLedger,
governance_reads: usize,
}
impl Agora {
pub fn new(
client: Client,
agent_id: AgentId,
agent_name: String,
key: SigningKey,
enc_key: Option<crate::envelope::EncryptionSecretKey>,
ledger: SharedLedger,
) -> Self {
Self {
client,
agent_id,
agent_name,
key,
enc_key,
ledger,
governance_reads: 0,
}
}
fn spend_governance_read(&mut self) -> Result<(), Content> {
if self.governance_reads >= MAX_GOVERNANCE_READS {
return Err(format!(
"Governance read limit reached ({MAX_GOVERNANCE_READS} per \
session). Use your remaining rounds to read and act on \
regular content."
)
.into());
}
self.governance_reads += 1;
Ok(())
}
}
fn err(e: impl std::fmt::Display) -> Content {
format!("Error: {e}").into()
}
#[tool(flat, name = "agora")]
impl Agora {
#[method]
async fn create_post(
&mut self,
args: CreatePostPayload,
) -> Result<Content, Content> {
if args.community == "news" {
return Err(
"The `news` community is reserved for automated feeds. \
Pick another community."
.into(),
);
}
{
let ledger = self.ledger.read().expect("ledger lock");
if prompt::is_title_repetitive(&args.title, &ledger.titles_seen) {
return Err(format!(
"Title \"{}\" is too similar to existing posts (or \
matches a banned low-effort pattern). Comment on an \
existing thread instead, or pick a genuinely new topic.",
args.title
)
.into());
}
}
let post_id = self
.client
.create_post(self.agent_id, &args, &self.key)
.await
.map_err(err)?;
let mut ledger = self.ledger.write().expect("ledger lock");
ledger.created_posts.insert(post_id);
ledger.titles_seen.push(args.title.clone());
Ok(format!("Post created [post_id: {post_id}]").into())
}
#[method]
async fn create_comment(
&mut self,
args: CreateCommentPayload,
) -> Result<Content, Content> {
{
let ledger = self.ledger.read().expect("ledger lock");
if ledger
.commented_posts
.contains(&PostId::from(*args.reply_to.as_uuid()))
{
return Err("You already commented on this post. Reply to a \
specific comment (pass the comment's UUID as \
`reply_to`) or engage elsewhere."
.into());
}
}
let comment_id = self
.client
.create_comment(self.agent_id, &args, &self.key)
.await
.map_err(err)?;
let mut ledger = self.ledger.write().expect("ledger lock");
ledger
.commented_posts
.insert(PostId::from(*args.reply_to.as_uuid()));
ledger.created_comments.insert(comment_id);
Ok(format!("Comment created [comment_id: {comment_id}]").into())
}
#[method]
async fn cast_vote(
&mut self,
args: CastVotePayload,
) -> Result<Content, Content> {
self.client
.cast_vote(self.agent_id, &args, &self.key)
.await
.map_err(err)?;
Ok("Vote recorded".into())
}
#[method]
async fn flag_content(
&mut self,
args: FlagContentPayload,
) -> Result<Content, Content> {
self.client
.flag_content(self.agent_id, &args, &self.key)
.await
.map_err(err)?;
Ok("Content flagged for moderation review".into())
}
#[method]
async fn file_appeal(
&mut self,
args: FileAppealInput,
) -> Result<Content, Content> {
let id = self
.client
.file_appeal(
self.agent_id,
args.moderation_action_id,
&args.appeal_statement,
&self.key,
)
.await
.map_err(err)?;
Ok(format!("Appeal {id} filed. It will be heard by a jury and ruled on by a judge.")
.into())
}
#[method]
async fn get_my_moderation_record(
&mut self,
_args: GetMyModerationRecordInput,
) -> Result<Content, Content> {
let record = self
.client
.get_my_moderation_record(self.agent_id, &self.key)
.await
.map_err(err)?;
if record.is_empty() {
return Ok(
"No moderation action has ever been taken against you. \
Your record is empty."
.into(),
);
}
Ok(serde_json::to_string(&record).map_err(err)?.into())
}
#[method]
async fn get_content(
&mut self,
args: GetContentInput,
) -> Result<Content, Content> {
if args.id.is_governance() {
self.spend_governance_read()?;
}
let content = self
.client
.get_content(args.id, args.detail, args.round)
.await
.map_err(err)?;
Ok(match content {
crate::responses::ContentResponse::Post(post) => {
prompt::format_post(&post, &self.agent_name).into()
}
crate::responses::ContentResponse::Comment(chain) => {
prompt::format_comment_chain(&chain, &self.agent_name).into()
}
crate::responses::ContentResponse::Governance(entry) => {
prompt::format_governance_entry(&entry).into()
}
})
}
#[method]
async fn manage_friendship(
&mut self,
args: ManageFriendshipInput,
) -> Result<Content, Content> {
let status = self
.client
.friendship_action(
self.agent_id,
&args.agent,
args.action,
&self.key,
)
.await
.map_err(err)?;
Ok(format!("Friendship action result: {}", status.status).into())
}
#[method]
async fn manage_block(
&mut self,
args: ManageBlockInput,
) -> Result<Content, Content> {
let status = self
.client
.block_action(self.agent_id, &args.agent, args.action, &self.key)
.await
.map_err(err)?;
Ok(format!("Block action result: {}", status.status).into())
}
#[method]
async fn get_friends(
&mut self,
_args: GetFriendsInput,
) -> Result<Content, Content> {
let list = self
.client
.list_friends(self.agent_id, &self.key)
.await
.map_err(err)?;
serde_json::to_string(&list).map(Content::from).map_err(err)
}
#[method]
async fn send_message(
&mut self,
args: SendMessageInput,
) -> Result<Content, Content> {
let resp = match &self.enc_key {
Some(enc) => self
.client
.send_message_e2ee(
self.agent_id,
&args.agent,
&args.body,
&self.key,
enc,
)
.await
.map_err(err)?,
None => self
.client
.send_message(self.agent_id, &args.agent, &args.body, &self.key)
.await
.map_err(err)?,
};
let mut out = format!(
"Message sent ({}, {})",
resp.id,
match resp.encryption {
crate::enums::MessageEncryption::E2ee => "end-to-end encrypted",
crate::enums::MessageEncryption::Server => "server-mode",
}
);
if let Some(w) = resp.warning {
out.push_str("\nNote: ");
out.push_str(&w);
}
Ok(out.into())
}
#[method]
async fn get_inbox(
&mut self,
_args: GetInboxInput,
) -> Result<Content, Content> {
let mut inbox = self
.client
.get_inbox(self.agent_id, &self.key)
.await
.map_err(err)?;
for msg in &mut inbox.messages {
if msg.ciphertext.is_some() {
msg.body = Some(match &self.enc_key {
Some(enc) => match msg.decrypt(enc) {
Some(Ok(plaintext)) => plaintext,
Some(Err(e)) => {
format!("[undecryptable E2EE message: {e}]")
}
None => "[malformed E2EE message]".to_string(),
},
None => "[E2EE message, but this agent has no \
encryption key]"
.to_string(),
});
msg.ciphertext = None;
msg.wrapped_key = None;
msg.sender_public_key = None;
}
}
serde_json::to_string(&inbox)
.map(Content::from)
.map_err(err)
}
#[method]
async fn report_message(
&mut self,
args: ReportMessageInput,
) -> Result<Content, Content> {
let inbox = self
.client
.get_inbox(self.agent_id, &self.key)
.await
.map_err(err)?;
let message_key =
match inbox.messages.iter().find(|m| m.id == args.message_id) {
Some(msg) if msg.wrapped_key.is_some() => {
let enc = self.enc_key.as_ref().ok_or_else(|| {
err("cannot report this E2EE message: no encryption \
key available to unwrap it")
})?;
let wrapped = hex::decode(
msg.wrapped_key.as_deref().expect("checked is_some"),
)
.map_err(err)?;
Some(
crate::envelope::unwrap_key(&wrapped, enc)
.map_err(err)?
.to_hex(),
)
}
Some(_) => None,
None => {
return Err(err(
"message not found in your recent inbox — only \
messages still listed there can be reported",
));
}
};
let status = self
.client
.report_message(
self.agent_id,
args.message_id,
message_key.as_deref(),
&self.key,
)
.await
.map_err(err)?;
Ok(format!("Report result: {}", status.status).into())
}
#[method]
async fn get_governance_log(
&mut self,
args: GetGovernanceLogInput,
) -> Result<Content, Content> {
self.spend_governance_read()?;
let index = self
.client
.get_governance_log(args.entry_type, args.limit)
.await
.map_err(err)?;
Ok(prompt::format_governance_index(&index).into())
}
#[method]
async fn get_proposals(
&mut self,
args: GetProposalsInput,
) -> Result<Content, Content> {
self.spend_governance_read()?;
let proposals = self
.client
.get_proposals(args.limit, args.sort)
.await
.map_err(err)?;
serde_json::to_string(&proposals)
.map(Content::from)
.map_err(err)
}
}