Skip to main content

telelog_server/archive/
mod.rs

1//! Long-term log storage in an object storage bucket you own (S3, R2, MinIO, GCS, or a local
2//! directory). The server archives every line it sees; the app queries it for history older
3//! than its in-memory buffer.
4
5mod chunk;
6pub mod ingest;
7
8use std::sync::{Arc, Mutex};
9use std::time::{Duration, SystemTime};
10
11use anyhow::{Context as _, Result, bail};
12use bytes::Bytes;
13use futures::{StreamExt, TryStreamExt};
14use object_store::aws::AmazonS3Builder;
15use object_store::gcp::GoogleCloudStorageBuilder;
16use object_store::local::LocalFileSystem;
17use object_store::memory::InMemory;
18use object_store::path::Path;
19use object_store::{ObjectStore, ObjectStoreExt, PutPayload};
20use serde::{Deserialize, Serialize};
21use telelog_core::{Filter, LogRecord};
22
23/// Most chunks fetched at once while answering a query.
24const FETCH_CONCURRENCY: usize = 8;
25/// Objects counted before storage stats give up and report a lower bound.
26const STATS_SCAN_LIMIT: u64 = 1_000_000;
27const CHECKPOINT: &str = "_checkpoint.json";
28
29/// What the ingest loop has done lately, for the Storage page.
30#[derive(Debug, Clone, Default)]
31pub struct Status {
32    pub last_flush: Option<SystemTime>,
33    pub last_error: Option<String>,
34    pub lines_archived: u64,
35}
36
37#[derive(Debug, Clone, Default)]
38pub struct Stats {
39    pub objects: u64,
40    pub bytes: u64,
41    /// Start of the oldest archived day.
42    pub oldest: Option<SystemTime>,
43    /// True when the bucket held more objects than were counted.
44    pub truncated: bool,
45}
46
47#[derive(Serialize, Deserialize)]
48struct Checkpoint {
49    /// Unix nanoseconds of the newest archived line.
50    last: i64,
51}
52
53pub struct Archive {
54    store: Arc<dyn ObjectStore>,
55    /// `<prefix>/v1`: everything this format writes lives under it.
56    root: Path,
57    /// The bucket URL without credentials, for display.
58    pub location: String,
59    pub provider: &'static str,
60    pub flush_every: Duration,
61    /// Days to keep; 0 keeps everything.
62    pub retention_days: u32,
63    status: Mutex<Status>,
64}
65
66impl Archive {
67    /// Opens `s3://bucket/prefix`, `gs://bucket/prefix`, `file:///path` or `memory://`.
68    /// S3 and GCS read credentials from the standard environment variables; for R2 or MinIO
69    /// set `AWS_ENDPOINT` (and `AWS_ALLOW_HTTP=true` for a plain-HTTP endpoint).
70    pub fn open(url: &str) -> Result<Self> {
71        let parsed = url::Url::parse(url).with_context(|| format!("invalid archive URL {url:?}"))?;
72        let host = parsed.host_str().unwrap_or_default();
73        let prefix = parsed.path().trim_matches('/');
74        let (store, provider, prefix): (Arc<dyn ObjectStore>, _, _) = match parsed.scheme() {
75            "s3" => {
76                let store = AmazonS3Builder::from_env()
77                    .with_url(format!("s3://{host}"))
78                    .build()
79                    .context("configuring S3")?;
80                let provider = if std::env::var_os("AWS_ENDPOINT").is_some() {
81                    "S3-compatible"
82                } else {
83                    "Amazon S3"
84                };
85                (Arc::new(store), provider, prefix)
86            }
87            "gs" => {
88                let store = GoogleCloudStorageBuilder::from_env()
89                    .with_url(format!("gs://{host}"))
90                    .build()
91                    .context("configuring Google Cloud Storage")?;
92                (Arc::new(store), "Google Cloud Storage", prefix)
93            }
94            "file" => {
95                let dir = parsed
96                    .to_file_path()
97                    .map_err(|_| anyhow::anyhow!("invalid file URL {url:?}"))?;
98                std::fs::create_dir_all(&dir).with_context(|| format!("creating {}", dir.display()))?;
99                (Arc::new(LocalFileSystem::new_with_prefix(&dir)?), "Local disk", "")
100            }
101            "memory" => (Arc::new(InMemory::new()), "In memory", ""),
102            other => bail!("unsupported archive URL scheme {other:?}; use s3://, gs://, file:// or memory://"),
103        };
104        let location = match parsed.scheme() {
105            "file" | "memory" => format!("{}://{}", parsed.scheme(), parsed.path()),
106            scheme if prefix.is_empty() => format!("{scheme}://{host}"),
107            scheme => format!("{scheme}://{host}/{prefix}"),
108        };
109        Ok(Self::with_store(store, prefix, location, provider))
110    }
111
112    pub fn with_store(store: Arc<dyn ObjectStore>, prefix: &str, location: String, provider: &'static str) -> Self {
113        let root = if prefix.is_empty() {
114            Path::from("v1")
115        } else {
116            Path::from(prefix).join("v1")
117        };
118        Archive {
119            store,
120            root,
121            location,
122            provider,
123            flush_every: Duration::from_secs(60),
124            retention_days: 0,
125            status: Mutex::default(),
126        }
127    }
128
129    pub fn status(&self) -> Status {
130        self.status.lock().unwrap().clone()
131    }
132
133    fn record_error(&self, error: &anyhow::Error) {
134        self.status.lock().unwrap().last_error = Some(format!("{error:#}"));
135    }
136
137    /// Writes one chunk and moves the checkpoint forward. `records` must not be empty.
138    pub async fn write_chunk(&self, records: &[LogRecord]) -> Result<()> {
139        let first = records.iter().map(|r| r.timestamp).min().context("empty chunk")?;
140        let last = records.iter().map(|r| r.timestamp).max().context("empty chunk")?;
141        let (date, hour) = chunk::partition(first);
142        let name = chunk::file_name(chunk::unix_millis(first), chunk::unix_millis(last), &random_id()?);
143        let path = self.root.clone().join(date).join(hour).join(name);
144        let body = chunk::encode(records)?;
145        self.store
146            .put(&path, PutPayload::from(Bytes::from(body)))
147            .await
148            .with_context(|| format!("writing {path}"))?;
149
150        let checkpoint = serde_json::to_vec(&Checkpoint {
151            last: chunk::unix_nanos(last),
152        })?;
153        self.store
154            .put(
155                &self.root.clone().join(CHECKPOINT),
156                PutPayload::from(Bytes::from(checkpoint)),
157            )
158            .await
159            .context("writing checkpoint")?;
160
161        let mut status = self.status.lock().unwrap();
162        status.last_flush = Some(SystemTime::now());
163        status.last_error = None;
164        status.lines_archived += records.len() as u64;
165        Ok(())
166    }
167
168    /// The newest archived line's time, if anything has been archived.
169    pub async fn checkpoint(&self) -> Result<Option<SystemTime>> {
170        match self.store.get(&self.root.clone().join(CHECKPOINT)).await {
171            Ok(object) => {
172                let checkpoint: Checkpoint = serde_json::from_slice(&object.bytes().await?)?;
173                Ok(Some(chunk::from_unix_nanos(checkpoint.last)))
174            }
175            Err(object_store::Error::NotFound { .. }) => Ok(None),
176            Err(e) => Err(e).context("reading checkpoint"),
177        }
178    }
179
180    /// Child "directories" of `dir`, e.g. the day or hour partitions.
181    async fn list_dirs(&self, dir: &Path) -> Result<Vec<String>> {
182        let listing = self.store.list_with_delimiter(Some(dir)).await?;
183        Ok(listing
184            .common_prefixes
185            .iter()
186            .filter_map(|p| p.filename().map(str::to_string))
187            .collect())
188    }
189
190    /// The newest `limit` archived lines matching `filter`, oldest first.
191    pub async fn query(&self, filter: &Filter, limit: usize) -> Result<Vec<LogRecord>> {
192        let to = filter.to.unwrap_or_else(SystemTime::now);
193        // Chunks live under the hour they start in and may run past it (flushes are at most
194        // an hour apart), so look one hour earlier than the range.
195        let from_key = filter.from.map(|from| {
196            let (date, hour) = chunk::partition(from.checked_sub(Duration::from_secs(3600)).unwrap_or(from));
197            format!("{date}/{hour}")
198        });
199        let (to_date, to_hour) = chunk::partition(to);
200        let to_key = format!("{to_date}/{to_hour}");
201        let (from_ms, to_ms) = (filter.from.map_or(i64::MIN, chunk::unix_millis), chunk::unix_millis(to));
202
203        let mut dates = self.list_dirs(&self.root).await?;
204        dates.retain(|d| d.as_str() <= to_date.as_str() && from_key.as_ref().is_none_or(|k| k[..10] <= *d.as_str()));
205        dates.sort_unstable_by(|a, b| b.cmp(a));
206
207        let mut found = Vec::new();
208        // Chunks within an hour aren't ordered, and one starting an hour earlier can still hold
209        // later lines, so keep reading one more hour after reaching the limit.
210        let mut hours_after_limit = 0;
211        'dates: for date in dates {
212            let date_dir = self.root.clone().join(date.as_str());
213            let mut hours = self.list_dirs(&date_dir).await?;
214            hours.sort_unstable_by(|a, b| b.cmp(a));
215            for hour in hours {
216                let key = format!("{date}/{hour}");
217                if key > to_key || from_key.as_ref().is_some_and(|k| key < *k) {
218                    continue;
219                }
220                let listing = self
221                    .store
222                    .list_with_delimiter(Some(&date_dir.clone().join(hour.as_str())))
223                    .await?;
224                let chunks: Vec<Path> = listing
225                    .objects
226                    .into_iter()
227                    .filter(|o| {
228                        o.location
229                            .filename()
230                            .and_then(chunk::parse_file_name)
231                            .is_some_and(|(start, end)| end >= from_ms && start <= to_ms)
232                    })
233                    .map(|o| o.location)
234                    .collect();
235                let mut fetches = futures::stream::iter(chunks)
236                    .map(|path| async move {
237                        let bytes = self.store.get(&path).await?.bytes().await?;
238                        chunk::decode(&bytes).with_context(|| format!("reading {path}"))
239                    })
240                    .buffer_unordered(FETCH_CONCURRENCY);
241                while let Some(records) = fetches.try_next().await? {
242                    found.extend(records.into_iter().filter(|r| filter.matches(r)));
243                }
244                if found.len() >= limit {
245                    hours_after_limit += 1;
246                    if hours_after_limit > 1 {
247                        break 'dates;
248                    }
249                }
250            }
251        }
252
253        found.sort_by_key(|r| r.timestamp);
254        if found.len() > limit {
255            found.drain(..found.len() - limit);
256        }
257        Ok(found)
258    }
259
260    /// Deletes whole days older than `keep_days` before `now`. Returns the objects removed.
261    pub async fn sweep(&self, keep_days: u32, now: SystemTime) -> Result<u64> {
262        let cutoff = now
263            .checked_sub(Duration::from_secs(u64::from(keep_days) * 86_400))
264            .unwrap_or(SystemTime::UNIX_EPOCH);
265        let (cutoff_date, _) = chunk::partition(cutoff);
266        let mut removed = 0;
267        for date in self.list_dirs(&self.root).await? {
268            if date >= cutoff_date {
269                continue;
270            }
271            let paths = self
272                .store
273                .list(Some(&self.root.clone().join(date.as_str())))
274                .map_ok(|o| o.location)
275                .boxed();
276            let mut deletions = self.store.delete_stream(paths);
277            while let Some(result) = deletions.next().await {
278                result?;
279                removed += 1;
280            }
281        }
282        Ok(removed)
283    }
284
285    pub async fn stats(&self) -> Result<Stats> {
286        let mut stats = Stats::default();
287        let mut objects = self.store.list(Some(&self.root));
288        while let Some(object) = objects.try_next().await? {
289            if object.location.filename() == Some(CHECKPOINT) {
290                continue;
291            }
292            if stats.objects == STATS_SCAN_LIMIT {
293                stats.truncated = true;
294                break;
295            }
296            stats.objects += 1;
297            stats.bytes += object.size;
298        }
299        let mut dates = self.list_dirs(&self.root).await?;
300        dates.sort_unstable();
301        stats.oldest = dates.first().and_then(|d| {
302            chrono::NaiveDate::parse_from_str(d, "%Y-%m-%d")
303                .ok()
304                .map(|d| d.and_time(chrono::NaiveTime::MIN).and_utc().into())
305        });
306        Ok(stats)
307    }
308}
309
310fn random_id() -> Result<String> {
311    let mut bytes = [0u8; 6];
312    getrandom::fill(&mut bytes).map_err(|e| anyhow::anyhow!("reading system randomness: {e}"))?;
313    Ok(bytes.iter().map(|b| format!("{b:02x}")).collect())
314}