Skip to main content

kcode_kennedy_session_objects/
lib.rs

1//! Object resolution and media staging for Kennedy sessions.
2
3#![forbid(unsafe_code)]
4
5use anyhow::Context as _;
6use kcode_kweb_db::ObjectId;
7use kcode_server_object_envelopes::{StoredFile, sanitize_file_name};
8use kcode_session_history::{
9    Session as HistorySession,
10    chatend::{ObjectMetadata, PendingId},
11};
12use kcode_telegram_session_coordinator::{Attachment, AttachmentRequest};
13use serde_json::{Value, json};
14
15/// Exact bytes and authoritative delivery metadata for a session object.
16#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct ResolvedObject {
18    pub object_id: String,
19    pub bytes: Vec<u8>,
20    pub file_name: String,
21    pub media_type: String,
22    pub transport_kind: Option<String>,
23}
24
25/// Authoritative facts for an object staged in the session journal.
26#[derive(Clone, Debug, Eq, PartialEq)]
27pub struct StagedDescriptor {
28    pub pending_id: String,
29    pub file_name: String,
30    pub media_type: String,
31    pub size_bytes: u64,
32    pub transport_kind: Option<String>,
33}
34
35impl StagedDescriptor {
36    /// Replaces caller-supplied descriptor facts with authoritative values.
37    pub fn apply_to(&self, descriptor: &mut Value) {
38        if !descriptor.is_object() {
39            *descriptor = json!({});
40        }
41        descriptor["pendingId"] = json!(self.pending_id);
42        descriptor["fileName"] = json!(self.file_name);
43        descriptor["extension"] = json!(file_name_extension(&self.file_name));
44        descriptor["mimeType"] = json!(self.media_type);
45        descriptor["sizeBytes"] = json!(self.size_bytes);
46    }
47}
48
49/// The result of staging or reusing one Telegram group media object.
50#[derive(Clone, Debug, Eq, PartialEq)]
51pub struct StagedTelegramMedia {
52    pub descriptor: StagedDescriptor,
53    pub kind: String,
54    pub reused: bool,
55}
56
57/// Correlation, limits, metadata, and timestamp for one Telegram staging call.
58#[derive(Clone, Debug, PartialEq)]
59pub struct TelegramStageRequest {
60    pub chat_id: i64,
61    pub message_id: i64,
62    pub maximum_bytes: u64,
63    pub transport_metadata: Value,
64    pub recorded_at: String,
65}
66
67/// Resolves a pending journal object or a canonical object through a caller
68/// supplied read boundary.
69pub fn resolve_object(
70    journal: &mut HistorySession,
71    object_id: &str,
72    read_canonical: impl FnOnce(&str) -> anyhow::Result<StoredFile>,
73) -> anyhow::Result<ResolvedObject> {
74    if object_id.starts_with("pending:") {
75        let pending_id = PendingId::parse(object_id.to_owned())?;
76        let location = journal
77            .objects()
78            .get(&pending_id)
79            .cloned()
80            .with_context(|| {
81                format!("staged object {pending_id} does not exist in this session")
82            })?;
83        let transport_kind = staged_object_transport_kind(journal, &pending_id);
84        let bytes = journal.read_object(&pending_id)?;
85        anyhow::ensure!(
86            bytes.len() as u64 == location.payload_len,
87            "staged object {pending_id} declared {} bytes but resolved to {}",
88            location.payload_len,
89            bytes.len()
90        );
91        let fallback = format!("object-{}.bin", pending_id.number());
92        return Ok(ResolvedObject {
93            object_id: pending_id.to_string(),
94            bytes,
95            file_name: sanitize_file_name(
96                location.metadata.file_name.as_deref().unwrap_or_default(),
97                &fallback,
98            ),
99            media_type: location.metadata.media_type,
100            transport_kind,
101        });
102    }
103
104    let canonical_id = object_id
105        .parse::<ObjectId>()
106        .with_context(|| format!("{object_id:?} is not an object ID"))?;
107    let file = read_canonical(object_id)?;
108    anyhow::ensure!(
109        file.object_id == canonical_id,
110        "object store returned {} while resolving {canonical_id}",
111        file.object_id
112    );
113    Ok(ResolvedObject {
114        object_id: canonical_id.to_string(),
115        bytes: file.bytes,
116        file_name: file.file_name,
117        media_type: file.media_type,
118        transport_kind: file.transport_kind,
119    })
120}
121
122/// Resolves and normalizes media while enforcing the current nonempty and size
123/// requirements.
124pub fn resolve_media_object(
125    journal: &mut HistorySession,
126    object_id: &str,
127    maximum_bytes: u64,
128    read_canonical: impl FnOnce(&str) -> anyhow::Result<StoredFile>,
129) -> anyhow::Result<ResolvedObject> {
130    let mut resolved = resolve_object(journal, object_id, read_canonical)?;
131    resolved.media_type = normalize_media_type(&resolved.media_type);
132    anyhow::ensure!(
133        !resolved.bytes.is_empty(),
134        "media object {} is empty",
135        resolved.object_id
136    );
137    anyhow::ensure!(
138        resolved.bytes.len() as u64 <= maximum_bytes,
139        "media object {} is {} bytes, over the {}-byte enrichment limit",
140        resolved.object_id,
141        resolved.bytes.len(),
142        maximum_bytes
143    );
144    Ok(resolved)
145}
146
147/// Returns authoritative staged metadata, using the same safe fallback filename
148/// as object resolution.
149pub fn staged_descriptor(
150    journal: &HistorySession,
151    pending_id: &PendingId,
152) -> anyhow::Result<StagedDescriptor> {
153    let location = journal
154        .objects()
155        .get(pending_id)
156        .with_context(|| format!("user-provided object {pending_id} is not staged"))?;
157    let fallback = format!("object-{}.bin", pending_id.number());
158    Ok(StagedDescriptor {
159        pending_id: pending_id.to_string(),
160        file_name: sanitize_file_name(
161            location.metadata.file_name.as_deref().unwrap_or_default(),
162            &fallback,
163        ),
164        media_type: normalize_media_type(&location.metadata.media_type),
165        size_bytes: location.payload_len,
166        transport_kind: staged_object_transport_kind(journal, pending_id),
167    })
168}
169
170/// Reuses media with matching Telegram correlation metadata or downloads and
171/// stages it once. The closures keep transport lookup and naming outside this
172/// crate.
173pub fn stage_telegram_group_media(
174    journal: &mut HistorySession,
175    request: TelegramStageRequest,
176    download: impl FnOnce() -> anyhow::Result<(Vec<u8>, String)>,
177    file_name_for_media_type: impl FnOnce(&str) -> String,
178) -> anyhow::Result<StagedTelegramMedia> {
179    if let Some((pending_id, metadata, size_bytes)) =
180        find_staged_telegram_group_media(journal, request.chat_id, request.message_id)
181    {
182        return telegram_stage_result(journal, pending_id, &metadata, size_bytes, true);
183    }
184
185    let (bytes, downloaded_media_type) = download()?;
186    anyhow::ensure!(
187        !bytes.is_empty(),
188        "Telegram group media message {} is empty",
189        request.message_id
190    );
191    anyhow::ensure!(
192        bytes.len() as u64 <= request.maximum_bytes,
193        "Telegram group media message {} is {} bytes, over the {}-byte enrichment limit",
194        request.message_id,
195        bytes.len(),
196        request.maximum_bytes
197    );
198    let media_type = normalize_media_type(&downloaded_media_type);
199    let file_name = file_name_for_media_type(&media_type);
200    let pending_id = journal.stage_object(
201        request.recorded_at,
202        media_type,
203        Some(file_name),
204        request.transport_metadata,
205        &bytes,
206    )?;
207    let metadata = journal
208        .objects()
209        .get(&pending_id)
210        .context("newly staged Telegram group media is missing")?
211        .metadata
212        .clone();
213    telegram_stage_result(journal, pending_id, &metadata, bytes.len() as u64, false)
214}
215
216/// Resolves delivery requests and constructs coordinator attachments without
217/// invoking a transport.
218pub fn delivery_attachments(
219    journal: &mut HistorySession,
220    requests: Vec<AttachmentRequest>,
221    mut read_canonical: impl FnMut(&str) -> anyhow::Result<StoredFile>,
222) -> anyhow::Result<Vec<Attachment>> {
223    requests
224        .into_iter()
225        .map(|request| {
226            let object = resolve_object(journal, &request.object_id, |id| read_canonical(id))?;
227            let file_name = request
228                .file_name
229                .unwrap_or_else(|| object.file_name.clone());
230            Ok(Attachment {
231                object_id: object.object_id,
232                bytes: object.bytes,
233                file_name,
234                media_type: object.media_type,
235                transport_kind: object.transport_kind,
236            })
237        })
238        .collect()
239}
240
241fn staged_object_transport_kind(
242    journal: &HistorySession,
243    pending_id: &PendingId,
244) -> Option<String> {
245    let pending_id_text = pending_id.to_string();
246    for state in journal.state().boxes.values() {
247        let Some(index) = state
248            .canonical
249            .content
250            .objects
251            .iter()
252            .position(|object_id| object_id == &pending_id_text)
253        else {
254            continue;
255        };
256        let metadata = &state.canonical.content.metadata;
257        let descriptor = metadata
258            .get("attachments")
259            .and_then(Value::as_array)
260            .and_then(|attachments| {
261                attachments
262                    .iter()
263                    .find(|attachment| {
264                        attachment.get("pendingId").and_then(Value::as_str)
265                            == Some(pending_id_text.as_str())
266                    })
267                    .or_else(|| attachments.get(index))
268            })
269            .or_else(|| metadata.get("media").filter(|value| value.is_object()));
270        if let Some(kind) = descriptor
271            .and_then(|descriptor| descriptor.get("kind"))
272            .and_then(Value::as_str)
273            .filter(|kind| !kind.trim().is_empty())
274        {
275            return Some(kind.to_owned());
276        }
277    }
278    journal
279        .objects()
280        .get(pending_id)
281        .and_then(|location| location.metadata.transport.get("kind"))
282        .and_then(Value::as_str)
283        .filter(|kind| !kind.trim().is_empty())
284        .map(str::to_owned)
285}
286
287fn find_staged_telegram_group_media(
288    journal: &HistorySession,
289    chat_id: i64,
290    message_id: i64,
291) -> Option<(PendingId, ObjectMetadata, u64)> {
292    journal.objects().iter().find_map(|(pending_id, location)| {
293        let transport = &location.metadata.transport;
294        (transport.get("source").and_then(Value::as_str) == Some("telegram-group")
295            && transport.get("chatId").and_then(Value::as_i64) == Some(chat_id)
296            && transport.get("messageId").and_then(Value::as_i64) == Some(message_id))
297        .then(|| {
298            (
299                pending_id.clone(),
300                location.metadata.clone(),
301                location.payload_len,
302            )
303        })
304    })
305}
306
307fn telegram_stage_result(
308    journal: &HistorySession,
309    pending_id: PendingId,
310    metadata: &ObjectMetadata,
311    size_bytes: u64,
312    reused: bool,
313) -> anyhow::Result<StagedTelegramMedia> {
314    let file_name = metadata
315        .file_name
316        .as_deref()
317        .filter(|value| !value.trim().is_empty())
318        .with_context(|| {
319            format!(
320                "staged object {} has no authoritative filename",
321                metadata.pending_id
322            )
323        })?
324        .to_owned();
325    Ok(StagedTelegramMedia {
326        descriptor: StagedDescriptor {
327            pending_id: pending_id.to_string(),
328            file_name,
329            media_type: normalize_media_type(&metadata.media_type),
330            size_bytes,
331            transport_kind: staged_object_transport_kind(journal, &pending_id),
332        },
333        kind: metadata
334            .transport
335            .get("kind")
336            .and_then(Value::as_str)
337            .unwrap_or("media")
338            .to_owned(),
339        reused,
340    })
341}
342
343fn normalize_media_type(value: &str) -> String {
344    value
345        .split(';')
346        .next()
347        .unwrap_or(value)
348        .trim()
349        .to_ascii_lowercase()
350}
351
352fn file_name_extension(file_name: &str) -> String {
353    file_name
354        .rsplit_once('.')
355        .and_then(|(stem, extension)| {
356            (!stem.is_empty() && !extension.is_empty()).then_some(extension)
357        })
358        .map(|extension| format!(".{extension}"))
359        .unwrap_or_else(|| "(none)".into())
360}
361
362#[cfg(test)]
363mod tests {
364    use super::*;
365
366    #[test]
367    fn descriptor_replaces_untrusted_file_facts() {
368        let descriptor = StagedDescriptor {
369            pending_id: "pending:7".into(),
370            file_name: "voice.OGG".into(),
371            media_type: "audio/ogg".into(),
372            size_bytes: 42,
373            transport_kind: Some("voice".into()),
374        };
375        let mut value = json!({
376            "pendingId":"pending:wrong",
377            "fileName":"../../wrong",
378            "mimeType":"text/plain",
379            "dataUrl":"retained only when caller has not sanitized it"
380        });
381        descriptor.apply_to(&mut value);
382        assert_eq!(value["pendingId"], "pending:7");
383        assert_eq!(value["fileName"], "voice.OGG");
384        assert_eq!(value["extension"], ".OGG");
385        assert_eq!(value["mimeType"], "audio/ogg");
386        assert_eq!(value["sizeBytes"], 42);
387    }
388
389    #[test]
390    fn normalization_is_parameter_insensitive() {
391        assert_eq!(
392            normalize_media_type(" Audio/OGG ; codecs=opus"),
393            "audio/ogg"
394        );
395        assert_eq!(file_name_extension(".hidden"), "(none)");
396        assert_eq!(file_name_extension("report.pdf"), ".pdf");
397    }
398}