1mod 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
23const FETCH_CONCURRENCY: usize = 8;
25const STATS_SCAN_LIMIT: u64 = 1_000_000;
27const CHECKPOINT: &str = "_checkpoint.json";
28
29#[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 pub oldest: Option<SystemTime>,
43 pub truncated: bool,
45}
46
47#[derive(Serialize, Deserialize)]
48struct Checkpoint {
49 last: i64,
51}
52
53pub struct Archive {
54 store: Arc<dyn ObjectStore>,
55 root: Path,
57 pub location: String,
59 pub provider: &'static str,
60 pub flush_every: Duration,
61 pub retention_days: u32,
63 status: Mutex<Status>,
64}
65
66impl Archive {
67 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 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 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 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 pub async fn query(&self, filter: &Filter, limit: usize) -> Result<Vec<LogRecord>> {
192 let to = filter.to.unwrap_or_else(SystemTime::now);
193 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 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 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}