feat: stability for deployment and distribution #44
11 changed files with 1736 additions and 48 deletions
13
Cargo.lock
generated
13
Cargo.lock
generated
|
|
@ -4647,6 +4647,19 @@ dependencies = [
|
|||
"ureq",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "swactor-node"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"clap",
|
||||
"ctrlc",
|
||||
"distribution",
|
||||
"iroh",
|
||||
"runtime-dashboard",
|
||||
"swactor",
|
||||
"swactor-datastore",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "swactor-std"
|
||||
version = "0.1.0"
|
||||
|
|
|
|||
|
|
@ -11,6 +11,7 @@ members = [
|
|||
"crates/datastore",
|
||||
"crates/shared-types",
|
||||
"crates/crypto-wasm",
|
||||
"crates/swactor-node",
|
||||
"tests/docker",
|
||||
"crates/ci",
|
||||
"crates/local-runner",
|
||||
|
|
|
|||
407
crates/datastore/src/bridge.rs
Normal file
407
crates/datastore/src/bridge.rs
Normal file
|
|
@ -0,0 +1,407 @@
|
|||
//! Bridge between the runtime dashboard's `DatastoreStatsProvider` trait and
|
||||
//! the datastore actor system. Allows the dashboard to perform CRUD operations
|
||||
//! and lifecycle management without depending on `swactor-datastore` types.
|
||||
|
||||
use std::collections::BTreeMap;
|
||||
use std::sync::Arc;
|
||||
use std::thread;
|
||||
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
||||
|
||||
use swactor::actor::ActorAddress;
|
||||
use swactor::runtime::{Inbox, Runtime};
|
||||
|
||||
use distribution::types::NodeId;
|
||||
use runtime_dashboard::datastore_collector::{
|
||||
DatastoreFactory, DatastoreStatsProvider, ListScope,
|
||||
};
|
||||
|
||||
use crate::actors::{BlobStoreActor, DatastoreNode, MetadataActor};
|
||||
use crate::chunking::reassemble_blob;
|
||||
use crate::messages::{DatastoreNodeMsg, DatastoreResponse, MetadataMsg};
|
||||
use crate::metrics::DatastoreMetrics;
|
||||
use crate::storage::{FilesystemBackend, InMemoryBackend};
|
||||
use crate::types::{ContentHash, DatastoreConfig};
|
||||
|
||||
const POLL_TIMEOUT: Duration = Duration::from_secs(5);
|
||||
const POLL_INTERVAL: Duration = Duration::from_millis(1);
|
||||
|
||||
fn poll_response(inbox: &Inbox<DatastoreResponse>, timeout: Duration) -> Option<DatastoreResponse> {
|
||||
let start = Instant::now();
|
||||
loop {
|
||||
if let Some(resp) = inbox.try_recv() {
|
||||
return Some(resp);
|
||||
}
|
||||
if start.elapsed() > timeout {
|
||||
return None;
|
||||
}
|
||||
thread::sleep(POLL_INTERVAL);
|
||||
}
|
||||
}
|
||||
|
||||
fn entry_to_json(entry: &crate::types::ObjectEntry) -> serde_json::Value {
|
||||
let node_hex: String = entry.node_id.0.iter().map(|b| format!("{b:02x}")).collect();
|
||||
serde_json::json!({
|
||||
"content_hash": entry.content_hash.to_hex(),
|
||||
"name": entry.name,
|
||||
"node_id": node_hex,
|
||||
"tags": entry.tags,
|
||||
"size_bytes": entry.size_bytes,
|
||||
"created_at": entry.created_at,
|
||||
})
|
||||
}
|
||||
|
||||
fn manifest_to_json(manifest: &crate::types::ObjectManifest) -> serde_json::Value {
|
||||
let chunks: Vec<serde_json::Value> = manifest
|
||||
.chunks
|
||||
.iter()
|
||||
.map(|c| {
|
||||
serde_json::json!({
|
||||
"hash": c.hash.to_hex(),
|
||||
"offset": c.offset,
|
||||
"size": c.size,
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
serde_json::json!({
|
||||
"content_hash": manifest.content_hash.to_hex(),
|
||||
"chunks": chunks,
|
||||
"total_size": manifest.total_size,
|
||||
"chunk_size": manifest.chunk_size,
|
||||
"content_type": manifest.content_type,
|
||||
})
|
||||
}
|
||||
|
||||
fn entries_to_json(entries: &[crate::types::ObjectEntry]) -> Vec<serde_json::Value> {
|
||||
entries.iter().map(entry_to_json).collect()
|
||||
}
|
||||
|
||||
/// Bridges the dashboard trait to the datastore actor system.
|
||||
pub struct DatastoreBridge {
|
||||
metrics: Arc<DatastoreMetrics>,
|
||||
runtime: Arc<Runtime>,
|
||||
datastore_addr: ActorAddress,
|
||||
metadata_addr: ActorAddress,
|
||||
#[allow(dead_code)]
|
||||
blob_store_addr: ActorAddress,
|
||||
}
|
||||
|
||||
impl DatastoreBridge {
|
||||
pub fn new(
|
||||
metrics: Arc<DatastoreMetrics>,
|
||||
runtime: Arc<Runtime>,
|
||||
datastore_addr: ActorAddress,
|
||||
metadata_addr: ActorAddress,
|
||||
blob_store_addr: ActorAddress,
|
||||
) -> Self {
|
||||
Self {
|
||||
metrics,
|
||||
runtime,
|
||||
datastore_addr,
|
||||
metadata_addr,
|
||||
blob_store_addr,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl DatastoreStatsProvider for DatastoreBridge {
|
||||
fn snapshot_json(&self) -> Option<String> {
|
||||
let snap = self.metrics.snapshot();
|
||||
serde_json::to_string(&snap).ok()
|
||||
}
|
||||
|
||||
fn is_running(&self) -> bool {
|
||||
true
|
||||
}
|
||||
|
||||
fn list_objects(&self, name_filter: Option<&str>, scope: ListScope) -> Result<String, String> {
|
||||
let inbox = self.runtime.new_inbox::<DatastoreResponse>()
|
||||
.map_err(|e| format!("failed to create inbox: {e}"))?;
|
||||
|
||||
match scope {
|
||||
ListScope::Local => {
|
||||
let _ = self.runtime.send_to(
|
||||
self.datastore_addr,
|
||||
DatastoreNodeMsg::List {
|
||||
name_filter: name_filter.map(|s| s.to_string()),
|
||||
all: false,
|
||||
reply_to: *inbox.addr(),
|
||||
},
|
||||
);
|
||||
}
|
||||
ListScope::Swarm => {
|
||||
let _ = self.runtime.send_to(
|
||||
self.metadata_addr,
|
||||
MetadataMsg::ListLocal {
|
||||
name_filter: name_filter.map(|s| s.to_string()),
|
||||
reply_to: *inbox.addr(),
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
match poll_response(&inbox, POLL_TIMEOUT) {
|
||||
Some(DatastoreResponse::ListOk { entries }) => {
|
||||
let json = serde_json::json!({ "entries": entries_to_json(&entries) }).to_string();
|
||||
Ok(json)
|
||||
}
|
||||
Some(DatastoreResponse::Error { reason }) => Err(reason),
|
||||
_ => Err("timeout".into()),
|
||||
}
|
||||
}
|
||||
|
||||
fn get_object(&self, hash: &str) -> Result<String, String> {
|
||||
let content_hash = ContentHash::from_hex(hash)
|
||||
.ok_or_else(|| "invalid content hash hex".to_string())?;
|
||||
|
||||
let inbox = self.runtime.new_inbox::<DatastoreResponse>()
|
||||
.map_err(|e| format!("failed to create inbox: {e}"))?;
|
||||
|
||||
let _ = self.runtime.send_to(
|
||||
self.datastore_addr,
|
||||
DatastoreNodeMsg::Get {
|
||||
content_hash,
|
||||
reply_to: *inbox.addr(),
|
||||
},
|
||||
);
|
||||
|
||||
match poll_response(&inbox, POLL_TIMEOUT) {
|
||||
Some(DatastoreResponse::GetOk { entry, manifest }) => {
|
||||
self.metrics.record_get(&content_hash.to_hex());
|
||||
let json = serde_json::json!({
|
||||
"entry": entry_to_json(&entry),
|
||||
"manifest": manifest_to_json(&manifest),
|
||||
})
|
||||
.to_string();
|
||||
Ok(json)
|
||||
}
|
||||
Some(DatastoreResponse::NotFound) => Err("not found".into()),
|
||||
Some(DatastoreResponse::Error { reason }) => Err(reason),
|
||||
_ => Err("timeout".into()),
|
||||
}
|
||||
}
|
||||
|
||||
fn get_data(&self, hash: &str) -> Result<Vec<u8>, String> {
|
||||
let content_hash = ContentHash::from_hex(hash)
|
||||
.ok_or_else(|| "invalid content hash hex".to_string())?;
|
||||
|
||||
// Get manifest
|
||||
let inbox = self.runtime.new_inbox::<DatastoreResponse>()
|
||||
.map_err(|e| format!("failed to create inbox: {e}"))?;
|
||||
|
||||
let _ = self.runtime.send_to(
|
||||
self.datastore_addr,
|
||||
DatastoreNodeMsg::Get {
|
||||
content_hash,
|
||||
reply_to: *inbox.addr(),
|
||||
},
|
||||
);
|
||||
|
||||
self.metrics.record_get(&content_hash.to_hex());
|
||||
|
||||
let manifest = match poll_response(&inbox, POLL_TIMEOUT) {
|
||||
Some(DatastoreResponse::GetOk { manifest, .. }) => manifest,
|
||||
Some(DatastoreResponse::NotFound) => return Err("not found".into()),
|
||||
Some(DatastoreResponse::Error { reason }) => return Err(reason),
|
||||
_ => return Err("timeout".into()),
|
||||
};
|
||||
|
||||
// Read chunks
|
||||
let mut chunk_data = Vec::new();
|
||||
for chunk_ref in &manifest.chunks {
|
||||
let chunk_inbox = self.runtime.new_inbox::<DatastoreResponse>()
|
||||
.map_err(|e| format!("failed to create inbox: {e}"))?;
|
||||
|
||||
let _ = self.runtime.send_to(
|
||||
self.datastore_addr,
|
||||
DatastoreNodeMsg::ReadChunk {
|
||||
hash: chunk_ref.hash,
|
||||
reply_to: *chunk_inbox.addr(),
|
||||
},
|
||||
);
|
||||
|
||||
match poll_response(&chunk_inbox, POLL_TIMEOUT) {
|
||||
Some(DatastoreResponse::ChunkOk { hash, data }) => {
|
||||
chunk_data.push((hash, data));
|
||||
}
|
||||
_ => return Err("failed to read chunk".into()),
|
||||
}
|
||||
}
|
||||
|
||||
reassemble_blob(&manifest, &chunk_data)
|
||||
.map_err(|e| format!("reassembly failed: {e:?}"))
|
||||
}
|
||||
|
||||
fn put_data(&self, data: Vec<u8>, name: Option<String>) -> Result<String, String> {
|
||||
let body_len = data.len();
|
||||
|
||||
let inbox = self.runtime.new_inbox::<DatastoreResponse>()
|
||||
.map_err(|e| format!("failed to create inbox: {e}"))?;
|
||||
|
||||
let _ = self.runtime.send_to(
|
||||
self.datastore_addr,
|
||||
DatastoreNodeMsg::Put {
|
||||
data,
|
||||
name: name.clone(),
|
||||
tags: BTreeMap::new(),
|
||||
reply_to: *inbox.addr(),
|
||||
},
|
||||
);
|
||||
|
||||
match poll_response(&inbox, POLL_TIMEOUT) {
|
||||
Some(DatastoreResponse::PutOk { content_hash }) => {
|
||||
let hex = content_hash.to_hex();
|
||||
self.metrics.record_put(&hex, name.as_deref(), body_len as u64);
|
||||
let json = serde_json::json!({ "content_hash": hex }).to_string();
|
||||
Ok(json)
|
||||
}
|
||||
Some(DatastoreResponse::Error { reason }) => Err(reason),
|
||||
_ => Err("timeout waiting for put response".into()),
|
||||
}
|
||||
}
|
||||
|
||||
fn delete_object(&self, hash: &str) -> Result<String, String> {
|
||||
let content_hash = ContentHash::from_hex(hash)
|
||||
.ok_or_else(|| "invalid content hash hex".to_string())?;
|
||||
|
||||
let inbox = self.runtime.new_inbox::<DatastoreResponse>()
|
||||
.map_err(|e| format!("failed to create inbox: {e}"))?;
|
||||
|
||||
let _ = self.runtime.send_to(
|
||||
self.datastore_addr,
|
||||
DatastoreNodeMsg::Delete {
|
||||
content_hash,
|
||||
reply_to: *inbox.addr(),
|
||||
},
|
||||
);
|
||||
|
||||
match poll_response(&inbox, POLL_TIMEOUT) {
|
||||
Some(DatastoreResponse::DeleteOk { content_hash }) => {
|
||||
let hex = content_hash.to_hex();
|
||||
self.metrics.record_delete(&hex, 0);
|
||||
let json = serde_json::json!({ "content_hash": hex }).to_string();
|
||||
Ok(json)
|
||||
}
|
||||
Some(DatastoreResponse::NotFound) => Err("not found".into()),
|
||||
Some(DatastoreResponse::Error { reason }) => Err(reason),
|
||||
_ => Err("timeout".into()),
|
||||
}
|
||||
}
|
||||
|
||||
fn node_status(&self) -> Result<String, String> {
|
||||
let inbox = self.runtime.new_inbox::<DatastoreResponse>()
|
||||
.map_err(|e| format!("failed to create inbox: {e}"))?;
|
||||
|
||||
let _ = self.runtime.send_to(
|
||||
self.datastore_addr,
|
||||
DatastoreNodeMsg::Status {
|
||||
reply_to: *inbox.addr(),
|
||||
},
|
||||
);
|
||||
|
||||
match poll_response(&inbox, POLL_TIMEOUT) {
|
||||
Some(DatastoreResponse::NodeStatus { node_id }) => {
|
||||
let hex: String = node_id.0.iter().map(|b| format!("{b:02x}")).collect();
|
||||
let json = serde_json::json!({ "node_id": hex }).to_string();
|
||||
Ok(json)
|
||||
}
|
||||
_ => Err("timeout".into()),
|
||||
}
|
||||
}
|
||||
|
||||
fn shutdown_datastore(&self) -> Result<(), String> {
|
||||
// We can't actually stop the actors from here without a runtime handle,
|
||||
// but we can signal shutdown. The caller (server handler) clears the
|
||||
// provider reference which effectively disables the datastore.
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
/// Factory that can spawn a new set of datastore actors on a shared runtime.
|
||||
pub struct DatastoreNodeFactory {
|
||||
runtime: Arc<Runtime>,
|
||||
default_chunk_size: u32,
|
||||
}
|
||||
|
||||
impl DatastoreNodeFactory {
|
||||
pub fn new(runtime: Arc<Runtime>, default_chunk_size: u32) -> Self {
|
||||
Self {
|
||||
runtime,
|
||||
default_chunk_size,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl DatastoreFactory for DatastoreNodeFactory {
|
||||
fn start_datastore(
|
||||
&self,
|
||||
storage_path: Option<String>,
|
||||
) -> Result<Arc<dyn DatastoreStatsProvider>, String> {
|
||||
// Generate a unique node ID
|
||||
let node_id = {
|
||||
let mut bytes = [0u8; 32];
|
||||
let nanos = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_nanos();
|
||||
for (i, b) in nanos.to_le_bytes().iter().enumerate() {
|
||||
bytes[i % 32] ^= *b;
|
||||
}
|
||||
let pid = std::process::id();
|
||||
for (i, b) in pid.to_le_bytes().iter().enumerate() {
|
||||
bytes[i + 16] ^= *b;
|
||||
}
|
||||
NodeId(bytes)
|
||||
};
|
||||
|
||||
let config = DatastoreConfig {
|
||||
chunk_size: self.default_chunk_size,
|
||||
storage_path: storage_path
|
||||
.as_ref()
|
||||
.map(|s| s.into())
|
||||
.unwrap_or_else(|| "datastore".into()),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let backend: Box<dyn crate::StorageBackend> = match &storage_path {
|
||||
Some(path) => {
|
||||
let p = std::path::PathBuf::from(path);
|
||||
std::fs::create_dir_all(&p)
|
||||
.map_err(|e| format!("failed to create storage directory: {e}"))?;
|
||||
Box::new(FilesystemBackend::new(p))
|
||||
}
|
||||
None => Box::new(InMemoryBackend::new()),
|
||||
};
|
||||
|
||||
let blob_store_addr = self
|
||||
.runtime
|
||||
.spawn(BlobStoreActor::new(backend))
|
||||
.map_err(|e| format!("failed to spawn BlobStoreActor: {e}"))?;
|
||||
|
||||
let mut metadata = MetadataActor::new(node_id, &config);
|
||||
metadata.set_blob_store(blob_store_addr);
|
||||
let metadata_addr = self
|
||||
.runtime
|
||||
.spawn(metadata)
|
||||
.map_err(|e| format!("failed to spawn MetadataActor: {e}"))?;
|
||||
|
||||
let datastore_node = DatastoreNode::new(node_id, blob_store_addr, metadata_addr, config);
|
||||
let datastore_addr = self
|
||||
.runtime
|
||||
.spawn(datastore_node)
|
||||
.map_err(|e| format!("failed to spawn DatastoreNode: {e}"))?;
|
||||
|
||||
let node_hex: String = node_id.0.iter().map(|b| format!("{b:02x}")).collect();
|
||||
let metrics = Arc::new(DatastoreMetrics::new());
|
||||
metrics.set_node_id(node_hex);
|
||||
|
||||
let bridge = DatastoreBridge::new(
|
||||
metrics,
|
||||
Arc::clone(&self.runtime),
|
||||
datastore_addr,
|
||||
metadata_addr,
|
||||
blob_store_addr,
|
||||
);
|
||||
|
||||
Ok(Arc::new(bridge))
|
||||
}
|
||||
}
|
||||
|
|
@ -10,6 +10,8 @@ pub mod metrics;
|
|||
pub mod api;
|
||||
#[cfg(feature = "node")]
|
||||
pub mod ui_html;
|
||||
#[cfg(feature = "node")]
|
||||
pub mod bridge;
|
||||
|
||||
pub use types::{ChunkRef, ContentHash, DatastoreConfig, ObjectEntry, ObjectManifest};
|
||||
pub use messages::{BlobStoreMsg, DatastoreNodeMsg, DatastoreResponse, MetadataMsg, TransferMsg};
|
||||
|
|
|
|||
|
|
@ -6,11 +6,70 @@
|
|||
//!
|
||||
//! The `swactor-datastore` crate implements this trait in its `node` feature.
|
||||
|
||||
/// Trait for providing datastore stats to the dashboard.
|
||||
use std::sync::Arc;
|
||||
|
||||
/// Scope filter for listing objects.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum ListScope {
|
||||
Local,
|
||||
Swarm,
|
||||
}
|
||||
|
||||
/// Trait for providing datastore stats and CRUD operations to the dashboard.
|
||||
///
|
||||
/// Implementations capture a point-in-time snapshot as serialized JSON.
|
||||
/// The dashboard polls this every ~200ms via SSE.
|
||||
///
|
||||
/// All command methods have default implementations returning `Err` so that
|
||||
/// existing `DatastoreMetrics` impls continue to compile without changes.
|
||||
pub trait DatastoreStatsProvider: Send + Sync {
|
||||
/// Return a JSON-serialized datastore snapshot, or `None` if unavailable.
|
||||
fn snapshot_json(&self) -> Option<String>;
|
||||
|
||||
/// List objects as a JSON string. `scope` selects local-only or swarm-wide.
|
||||
fn list_objects(&self, _name_filter: Option<&str>, _scope: ListScope) -> Result<String, String> {
|
||||
Err("not supported".into())
|
||||
}
|
||||
|
||||
/// Get a single object's metadata + manifest as JSON.
|
||||
fn get_object(&self, _hash: &str) -> Result<String, String> {
|
||||
Err("not supported".into())
|
||||
}
|
||||
|
||||
/// Get the raw binary data for an object.
|
||||
fn get_data(&self, _hash: &str) -> Result<Vec<u8>, String> {
|
||||
Err("not supported".into())
|
||||
}
|
||||
|
||||
/// Store data, optionally with a name. Returns JSON with `content_hash`.
|
||||
fn put_data(&self, _data: Vec<u8>, _name: Option<String>) -> Result<String, String> {
|
||||
Err("not supported".into())
|
||||
}
|
||||
|
||||
/// Delete an object by hash. Returns JSON confirmation.
|
||||
fn delete_object(&self, _hash: &str) -> Result<String, String> {
|
||||
Err("not supported".into())
|
||||
}
|
||||
|
||||
/// Get node status as JSON.
|
||||
fn node_status(&self) -> Result<String, String> {
|
||||
Err("not supported".into())
|
||||
}
|
||||
|
||||
/// Whether the datastore is currently running.
|
||||
fn is_running(&self) -> bool {
|
||||
false
|
||||
}
|
||||
|
||||
/// Shut down the datastore actors.
|
||||
fn shutdown_datastore(&self) -> Result<(), String> {
|
||||
Err("not supported".into())
|
||||
}
|
||||
}
|
||||
|
||||
/// Factory for creating a new datastore instance from the dashboard.
|
||||
pub trait DatastoreFactory: Send + Sync {
|
||||
/// Start a datastore with optional persistent storage path.
|
||||
/// Returns a provider that can be installed into the dashboard.
|
||||
fn start_datastore(&self, storage_path: Option<String>) -> Result<Arc<dyn DatastoreStatsProvider>, String>;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -80,6 +80,8 @@ pub const DATASTORE_HTML: &str = r##"<!DOCTYPE html>
|
|||
white-space: nowrap;
|
||||
}
|
||||
.objects-table th { color: #888; font-weight: 500; position: sticky; top: 0; background: #161822; }
|
||||
.objects-table tr { cursor: pointer; }
|
||||
.objects-table tr:hover td { background: #1c1f2e; }
|
||||
|
||||
.transfers-section { display: none; }
|
||||
.transfers-section.visible { display: block; }
|
||||
|
|
@ -96,6 +98,87 @@ pub const DATASTORE_HTML: &str = r##"<!DOCTYPE html>
|
|||
}
|
||||
.transfer-label { color: #888; font-size: 11px; min-width: 80px; text-align: right; }
|
||||
|
||||
/* Upload panel */
|
||||
.upload-row {
|
||||
display: flex; gap: 10px; align-items: center; flex-wrap: wrap;
|
||||
}
|
||||
input[type="file"] {
|
||||
background: #1c1f2e; color: #e0e0e0; border: 1px solid #2a2d3e;
|
||||
border-radius: 4px; padding: 8px; font-family: inherit; font-size: 13px;
|
||||
min-height: 38px; cursor: pointer;
|
||||
}
|
||||
input[type="file"]::file-selector-button {
|
||||
background: #1e2030; color: #e0e0e0; border: 1px solid #2a2d3e;
|
||||
border-radius: 4px; padding: 4px 10px; font-family: inherit;
|
||||
font-size: 12px; cursor: pointer; margin-right: 8px;
|
||||
}
|
||||
input[type="file"]::file-selector-button:hover { border-color: #6366f1; }
|
||||
input[type="text"] {
|
||||
background: #1c1f2e; color: #e0e0e0; border: 1px solid #2a2d3e;
|
||||
border-radius: 4px; padding: 6px 10px; font-family: inherit;
|
||||
font-size: 13px; min-height: 38px; width: 180px;
|
||||
}
|
||||
input[type="text"]:focus { outline: none; border-color: #6366f1; }
|
||||
|
||||
/* Buttons */
|
||||
button {
|
||||
background: #1e2030; color: #e0e0e0; border: 1px solid #2a2d3e;
|
||||
border-radius: 4px; padding: 6px 14px; font-family: inherit;
|
||||
font-size: 13px; cursor: pointer; min-height: 38px;
|
||||
transition: border-color 0.15s;
|
||||
}
|
||||
button:hover { border-color: #6366f1; color: #fff; }
|
||||
button:disabled { opacity: 0.4; cursor: default; }
|
||||
button.danger:hover { border-color: #f44336; }
|
||||
button.primary { background: #6366f1; border-color: #6366f1; color: #fff; font-weight: 600; }
|
||||
button.primary:hover { background: #5558e6; }
|
||||
|
||||
/* Actions in table */
|
||||
.actions-cell { white-space: nowrap; text-align: right; }
|
||||
.actions-cell button { min-height: 28px; padding: 2px 8px; font-size: 11px; }
|
||||
|
||||
/* Origin badge */
|
||||
.origin-badge {
|
||||
display: inline-block; font-size: 10px; padding: 1px 6px;
|
||||
border-radius: 3px; font-weight: 600;
|
||||
}
|
||||
.origin-badge.local { background: #1b3a2a; color: #4caf50; }
|
||||
.origin-badge.remote { background: #1a2a3e; color: #2196f3; }
|
||||
|
||||
/* Toast */
|
||||
.toast {
|
||||
position: fixed; bottom: 20px; right: 20px; padding: 10px 16px;
|
||||
border-radius: 4px; font-size: 12px; z-index: 100; opacity: 0;
|
||||
transition: opacity 0.3s; pointer-events: none;
|
||||
}
|
||||
.toast.show { opacity: 1; }
|
||||
.toast.success { background: #4caf50; color: #fff; }
|
||||
.toast.error { background: #f44336; color: #fff; }
|
||||
|
||||
/* Modal */
|
||||
.modal-overlay {
|
||||
display: none; position: fixed; inset: 0;
|
||||
background: rgba(0,0,0,0.6); z-index: 50;
|
||||
align-items: center; justify-content: center;
|
||||
}
|
||||
.modal-overlay.open { display: flex; }
|
||||
.modal {
|
||||
background: #161822; border: 1px solid #2a2d3e; border-radius: 6px;
|
||||
padding: 20px; width: 90%; max-width: 560px; max-height: 80vh;
|
||||
overflow-y: auto;
|
||||
}
|
||||
.modal h2 { font-size: 14px; color: #fff; margin-bottom: 16px; text-transform: none; letter-spacing: 0; }
|
||||
.modal-close {
|
||||
float: right; background: none; border: none; color: #888;
|
||||
font-size: 18px; cursor: pointer; min-height: auto; padding: 0;
|
||||
}
|
||||
.modal-close:hover { color: #fff; border: none; }
|
||||
.detail-row { display: flex; margin-bottom: 8px; }
|
||||
.detail-label { color: #888; width: 110px; flex-shrink: 0; font-size: 11px; text-transform: uppercase; padding-top: 2px; }
|
||||
.detail-value { color: #e0e0e0; word-break: break-all; font-size: 13px; }
|
||||
.chunk-list { margin-top: 8px; }
|
||||
.chunk-item { color: #888; font-size: 11px; padding: 2px 0; }
|
||||
|
||||
::-webkit-scrollbar { width: 6px; }
|
||||
::-webkit-scrollbar-track { background: #0f1117; }
|
||||
::-webkit-scrollbar-thumb { background: #2a2d3e; border-radius: 3px; }
|
||||
|
|
@ -117,10 +200,22 @@ pub const DATASTORE_HTML: &str = r##"<!DOCTYPE html>
|
|||
</div>
|
||||
<div class="header-right">
|
||||
<span id="nodeLabel" style="color:#888;font-size:12px;">Waiting for data...</span>
|
||||
<button id="startBtn" onclick="startDatastore()" style="display:none" class="primary">Start Datastore</button>
|
||||
<button id="stopBtn" onclick="stopDatastore()" style="display:none" class="danger">Stop</button>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<div class="grid">
|
||||
<!-- Upload panel -->
|
||||
<div class="panel full-width" id="uploadPanel" style="display:none">
|
||||
<h2>Upload</h2>
|
||||
<div class="upload-row">
|
||||
<input type="file" id="fileInput" />
|
||||
<input type="text" id="nameInput" placeholder="name (optional)" />
|
||||
<button class="primary" id="uploadBtn" onclick="upload()">Upload</button>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<!-- Stat cards -->
|
||||
<div class="panel full-width">
|
||||
<h2>Datastore Stats</h2>
|
||||
|
|
@ -142,9 +237,9 @@ pub const DATASTORE_HTML: &str = r##"<!DOCTYPE html>
|
|||
<!-- Objects table -->
|
||||
<div class="panel">
|
||||
<h2>Objects <span id="objectCount" style="color:#555;font-weight:400;"></span></h2>
|
||||
<div class="objects-table">
|
||||
<div class="objects-table" id="objectsTableWrap">
|
||||
<table>
|
||||
<thead><tr><th>Hash</th><th>Name</th><th>Size</th></tr></thead>
|
||||
<thead><tr><th>Hash</th><th>Name</th><th>Origin</th><th>Size</th><th style="text-align:right">Actions</th></tr></thead>
|
||||
<tbody id="objectsBody"></tbody>
|
||||
</table>
|
||||
</div>
|
||||
|
|
@ -157,10 +252,25 @@ pub const DATASTORE_HTML: &str = r##"<!DOCTYPE html>
|
|||
</div>
|
||||
</div>
|
||||
|
||||
<!-- Detail modal -->
|
||||
<div class="modal-overlay" id="modal" onclick="if(event.target===this)closeModal()">
|
||||
<div class="modal">
|
||||
<button class="modal-close" onclick="closeModal()">×</button>
|
||||
<h2>Object Detail</h2>
|
||||
<div id="modalBody"></div>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<div class="toast" id="toast"></div>
|
||||
|
||||
<script>
|
||||
(function() {
|
||||
var DASHBOARD_MODE = '__DASHBOARD_MODE__';
|
||||
var dot = document.getElementById('statusDot');
|
||||
var dsRunning = false;
|
||||
var currentNodeId = '';
|
||||
|
||||
function $(id) { return document.getElementById(id); }
|
||||
|
||||
function formatBytes(b) {
|
||||
if (b === 0) return '0 B';
|
||||
|
|
@ -182,24 +292,59 @@ pub const DATASTORE_HTML: &str = r##"<!DOCTYPE html>
|
|||
return s.replace(/&/g,'&').replace(/</g,'<').replace(/>/g,'>');
|
||||
}
|
||||
|
||||
function updateFromSnapshot(snap) {
|
||||
function escAttr(s) {
|
||||
if (!s) return '';
|
||||
return s.replace(/\\/g,'\\\\').replace(/'/g,"\\'").replace(/"/g,'"');
|
||||
}
|
||||
|
||||
// Toast notification
|
||||
window.toast = function(msg, type) {
|
||||
var t = $('toast');
|
||||
t.textContent = msg;
|
||||
t.className = 'toast show ' + type;
|
||||
setTimeout(function() { t.className = 'toast'; }, 2500);
|
||||
};
|
||||
|
||||
function detailRow(label, value) {
|
||||
return '<div class="detail-row"><div class="detail-label">' + label +
|
||||
'</div><div class="detail-value">' + escapeHtml(String(value)) + '</div></div>';
|
||||
}
|
||||
|
||||
// Update lifecycle buttons
|
||||
function updateLifecycleUI(running) {
|
||||
dsRunning = running;
|
||||
$('startBtn').style.display = running ? 'none' : 'inline-block';
|
||||
$('stopBtn').style.display = running ? 'inline-block' : 'none';
|
||||
$('uploadPanel').style.display = running ? 'block' : 'none';
|
||||
}
|
||||
|
||||
function updateFromSnapshot(data) {
|
||||
// Handle envelope format: {is_running, snapshot}
|
||||
var running = data.is_running;
|
||||
var snap = data.snapshot;
|
||||
|
||||
updateLifecycleUI(running);
|
||||
|
||||
if (!snap) return;
|
||||
|
||||
// Node label
|
||||
if (snap.node_id) {
|
||||
document.getElementById('nodeLabel').textContent = 'Node: ' + snap.node_id.substring(0, 16) + '\u2026';
|
||||
currentNodeId = snap.node_id;
|
||||
$('nodeLabel').textContent = 'Node: ' + snap.node_id.substring(0, 16) + '\u2026';
|
||||
}
|
||||
|
||||
// Stat cards
|
||||
document.getElementById('statObjects').textContent = snap.object_count;
|
||||
document.getElementById('statSize').textContent = formatBytes(snap.total_bytes);
|
||||
document.getElementById('statPuts').textContent = snap.put_ops;
|
||||
document.getElementById('statGets').textContent = snap.get_ops;
|
||||
document.getElementById('statDeletes').textContent = snap.delete_ops;
|
||||
$('statObjects').textContent = snap.object_count;
|
||||
$('statSize').textContent = formatBytes(snap.total_bytes);
|
||||
$('statPuts').textContent = snap.put_ops;
|
||||
$('statGets').textContent = snap.get_ops;
|
||||
$('statDeletes').textContent = snap.delete_ops;
|
||||
|
||||
// Event timeline
|
||||
var timeline = document.getElementById('eventTimeline');
|
||||
var timeline = $('eventTimeline');
|
||||
var wasAtBottom = timeline.scrollTop + timeline.clientHeight >= timeline.scrollHeight - 20;
|
||||
timeline.innerHTML = '';
|
||||
document.getElementById('eventCount').textContent = '(' + snap.recent_events.length + ')';
|
||||
$('eventCount').textContent = '(' + snap.recent_events.length + ')';
|
||||
|
||||
for (var i = snap.recent_events.length - 1; i >= 0; i--) {
|
||||
var ev = snap.recent_events[i];
|
||||
|
|
@ -219,23 +364,36 @@ pub const DATASTORE_HTML: &str = r##"<!DOCTYPE html>
|
|||
}
|
||||
|
||||
// Objects table
|
||||
var tbody = document.getElementById('objectsBody');
|
||||
var tbody = $('objectsBody');
|
||||
tbody.innerHTML = '';
|
||||
document.getElementById('objectCount').textContent = '(' + snap.objects.length + ')';
|
||||
$('objectCount').textContent = '(' + snap.objects.length + ')';
|
||||
|
||||
for (var i = 0; i < snap.objects.length; i++) {
|
||||
var obj = snap.objects[i];
|
||||
var tr = document.createElement('tr');
|
||||
var nodeHex = obj.node_id || currentNodeId || '';
|
||||
var isLocal = !obj.node_id || obj.node_id === currentNodeId;
|
||||
var badgeClass = isLocal ? 'local' : 'remote';
|
||||
var badgeText = isLocal ? 'local' : (nodeHex ? nodeHex.substring(0, 8) : 'remote');
|
||||
var h = obj.hash || obj.content_hash || '';
|
||||
var shortHash = h.substring(0, 16);
|
||||
var name = obj.name || '\u2014';
|
||||
tr.setAttribute('onclick', "showDetail('" + escAttr(h) + "')");
|
||||
tr.innerHTML =
|
||||
'<td style="color:#aaa;font-size:10px;">' + obj.hash.substring(0, 16) + '\u2026</td>' +
|
||||
'<td>' + escapeHtml(obj.name || '\u2014') + '</td>' +
|
||||
'<td style="color:#888;">' + formatBytes(obj.size_bytes) + '</td>';
|
||||
'<td style="color:#6366f1;font-size:11px;" title="' + escapeHtml(h) + '">' + shortHash + '\u2026</td>' +
|
||||
'<td>' + escapeHtml(name) + '</td>' +
|
||||
'<td><span class="origin-badge ' + badgeClass + '">' + escapeHtml(badgeText) + '</span></td>' +
|
||||
'<td style="color:#888;">' + formatBytes(obj.size_bytes) + '</td>' +
|
||||
'<td class="actions-cell">' +
|
||||
'<button onclick="event.stopPropagation();download(\'' + escAttr(h) + '\',\'' + escAttr(obj.name || shortHash) + '\')">download</button> ' +
|
||||
'<button class="danger" onclick="event.stopPropagation();del(\'' + escAttr(h) + '\')">delete</button>' +
|
||||
'</td>';
|
||||
tbody.appendChild(tr);
|
||||
}
|
||||
|
||||
// Transfers
|
||||
var panel = document.getElementById('transfersPanel');
|
||||
var list = document.getElementById('transfersList');
|
||||
var panel = $('transfersPanel');
|
||||
var list = $('transfersList');
|
||||
|
||||
if (snap.active_transfers.length === 0) {
|
||||
panel.className = 'panel full-width transfers-section';
|
||||
|
|
@ -246,24 +404,151 @@ pub const DATASTORE_HTML: &str = r##"<!DOCTYPE html>
|
|||
for (var i = 0; i < snap.active_transfers.length; i++) {
|
||||
var t = snap.active_transfers[i];
|
||||
var pct = t.chunks_total > 0 ? Math.round((t.chunks_received / t.chunks_total) * 100) : 0;
|
||||
var row = document.createElement('div');
|
||||
row.className = 'transfer-row';
|
||||
row.innerHTML =
|
||||
var trow = document.createElement('div');
|
||||
trow.className = 'transfer-row';
|
||||
trow.innerHTML =
|
||||
'<span class="transfer-hash">' + t.hash.substring(0, 16) + '\u2026</span>' +
|
||||
'<div class="transfer-bar"><div class="transfer-fill" style="width:' + pct + '%"></div></div>' +
|
||||
'<span class="transfer-label">' + t.chunks_received + ' / ' + t.chunks_total + '</span>';
|
||||
list.appendChild(row);
|
||||
list.appendChild(trow);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SSE connection
|
||||
// ── CRUD operations ──────────────────────────────────────────────────
|
||||
|
||||
window.upload = function() {
|
||||
var file = $('fileInput').files[0];
|
||||
if (!file) { toast('select a file first', 'error'); return; }
|
||||
var name = $('nameInput').value.trim();
|
||||
var btn = $('uploadBtn');
|
||||
btn.disabled = true;
|
||||
btn.textContent = 'uploading...';
|
||||
var url = '/api/datastore/put';
|
||||
if (name) url += '?name=' + encodeURIComponent(name);
|
||||
fetch(url, { method: 'POST', body: file })
|
||||
.then(function(r) {
|
||||
if (!r.ok) return r.json().then(function(j) { throw new Error(j.error || r.statusText); });
|
||||
return r.json();
|
||||
})
|
||||
.then(function(j) {
|
||||
toast('uploaded ' + j.content_hash.substring(0, 12), 'success');
|
||||
$('fileInput').value = '';
|
||||
$('nameInput').value = '';
|
||||
})
|
||||
.catch(function(e) { toast('upload failed: ' + e.message, 'error'); })
|
||||
.finally(function() { btn.disabled = false; btn.textContent = 'Upload'; });
|
||||
};
|
||||
|
||||
window.download = function(hash, filename) {
|
||||
fetch('/api/datastore/data?hash=' + hash)
|
||||
.then(function(r) {
|
||||
if (!r.ok) throw new Error('not found');
|
||||
return r.blob();
|
||||
})
|
||||
.then(function(blob) {
|
||||
var a = document.createElement('a');
|
||||
a.href = URL.createObjectURL(blob);
|
||||
a.download = filename;
|
||||
a.click();
|
||||
URL.revokeObjectURL(a.href);
|
||||
})
|
||||
.catch(function(e) { toast('download failed: ' + e.message, 'error'); });
|
||||
};
|
||||
|
||||
window.del = function(hash) {
|
||||
if (!confirm('Delete ' + hash.substring(0, 16) + '?')) return;
|
||||
fetch('/api/datastore/delete?hash=' + hash, { method: 'POST' })
|
||||
.then(function(r) {
|
||||
if (!r.ok) return r.json().then(function(j) { throw new Error(j.error || r.statusText); });
|
||||
toast('deleted', 'success');
|
||||
})
|
||||
.catch(function(e) { toast('delete failed: ' + e.message, 'error'); });
|
||||
};
|
||||
|
||||
window.showDetail = function(hash) {
|
||||
fetch('/api/datastore/get?hash=' + hash)
|
||||
.then(function(r) {
|
||||
if (!r.ok) throw new Error('not found');
|
||||
return r.json();
|
||||
})
|
||||
.then(function(j) {
|
||||
var e = j.entry;
|
||||
var m = j.manifest;
|
||||
var html = '';
|
||||
html += detailRow('Hash', e.content_hash);
|
||||
html += detailRow('Name', e.name || '\u2014');
|
||||
html += detailRow('Size', formatBytes(e.size_bytes));
|
||||
html += detailRow('Node', e.node_id ? e.node_id.substring(0, 16) + '...' : '\u2014');
|
||||
if (e.tags && Object.keys(e.tags).length > 0) {
|
||||
html += detailRow('Tags', Object.entries(e.tags).map(function(kv) { return kv[0] + '=' + kv[1]; }).join(', '));
|
||||
}
|
||||
html += detailRow('Chunks', m.chunks.length + ' (' + formatBytes(m.chunk_size) + ' each)');
|
||||
if (m.chunks.length > 0) {
|
||||
html += '<div class="chunk-list">';
|
||||
for (var i = 0; i < m.chunks.length; i++) {
|
||||
var c = m.chunks[i];
|
||||
html += '<div class="chunk-item">#' + i + ' ' + c.hash.substring(0, 16) + ' (' + formatBytes(c.size) + ')</div>';
|
||||
}
|
||||
html += '</div>';
|
||||
}
|
||||
$('modalBody').innerHTML = html;
|
||||
$('modal').classList.add('open');
|
||||
})
|
||||
.catch(function(e) { toast('failed to load detail', 'error'); });
|
||||
};
|
||||
|
||||
window.closeModal = function() { $('modal').classList.remove('open'); };
|
||||
document.addEventListener('keydown', function(e) { if (e.key === 'Escape') closeModal(); });
|
||||
|
||||
// ── Lifecycle ────────────────────────────────────────────────────────
|
||||
|
||||
window.startDatastore = function() {
|
||||
$('startBtn').disabled = true;
|
||||
fetch('/api/datastore/start', { method: 'POST' })
|
||||
.then(function(r) {
|
||||
if (!r.ok) return r.json().then(function(j) { throw new Error(j.error || r.statusText); });
|
||||
return r.json();
|
||||
})
|
||||
.then(function() {
|
||||
toast('datastore started', 'success');
|
||||
updateLifecycleUI(true);
|
||||
})
|
||||
.catch(function(e) { toast('start failed: ' + e.message, 'error'); })
|
||||
.finally(function() { $('startBtn').disabled = false; });
|
||||
};
|
||||
|
||||
window.stopDatastore = function() {
|
||||
if (!confirm('Stop the datastore? In-memory data will be lost.')) return;
|
||||
$('stopBtn').disabled = true;
|
||||
fetch('/api/datastore/shutdown', { method: 'POST' })
|
||||
.then(function(r) {
|
||||
if (!r.ok) return r.json().then(function(j) { throw new Error(j.error || r.statusText); });
|
||||
return r.json();
|
||||
})
|
||||
.then(function() {
|
||||
toast('datastore stopped', 'success');
|
||||
updateLifecycleUI(false);
|
||||
$('objectsBody').innerHTML = '';
|
||||
$('eventTimeline').innerHTML = '';
|
||||
$('statObjects').textContent = '0';
|
||||
$('statSize').textContent = '0 B';
|
||||
$('statPuts').textContent = '0';
|
||||
$('statGets').textContent = '0';
|
||||
$('statDeletes').textContent = '0';
|
||||
})
|
||||
.catch(function(e) { toast('stop failed: ' + e.message, 'error'); })
|
||||
.finally(function() { $('stopBtn').disabled = false; });
|
||||
};
|
||||
|
||||
// ── SSE connection ───────────────────────────────────────────────────
|
||||
|
||||
var es = new EventSource('/events');
|
||||
|
||||
es.addEventListener('datastore', function(e) {
|
||||
try {
|
||||
var snap = JSON.parse(e.data);
|
||||
updateFromSnapshot(snap);
|
||||
var data = JSON.parse(e.data);
|
||||
updateFromSnapshot(data);
|
||||
} catch(err) { console.error('datastore parse error', err); }
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -98,6 +98,7 @@ pub struct DashboardHandle {
|
|||
#[cfg(feature = "distribution")]
|
||||
distribution: Arc<Mutex<Option<Arc<dyn distribution_collector::DistributionStatsProvider>>>>,
|
||||
datastore: Arc<Mutex<Option<Arc<dyn datastore_collector::DatastoreStatsProvider>>>>,
|
||||
datastore_factory: Arc<Mutex<Option<Arc<dyn datastore_collector::DatastoreFactory>>>>,
|
||||
#[cfg(feature = "ci")]
|
||||
ci: Arc<Mutex<Option<Arc<dyn ci_collector::CiStatsProvider>>>>,
|
||||
}
|
||||
|
|
@ -137,6 +138,16 @@ impl DashboardHandle {
|
|||
*self.datastore.lock().unwrap() = Some(provider);
|
||||
}
|
||||
|
||||
/// Attach a datastore factory, enabling start/stop from the dashboard.
|
||||
pub fn set_datastore_factory(&self, factory: Arc<dyn datastore_collector::DatastoreFactory>) {
|
||||
*self.datastore_factory.lock().unwrap() = Some(factory);
|
||||
}
|
||||
|
||||
/// Get the shared datastore provider mutex (for external wiring).
|
||||
pub fn datastore_provider(&self) -> &Arc<Mutex<Option<Arc<dyn datastore_collector::DatastoreStatsProvider>>>> {
|
||||
&self.datastore
|
||||
}
|
||||
|
||||
/// Attach a CI stats provider, enabling the `/api/ci/*` endpoints.
|
||||
#[cfg(feature = "ci")]
|
||||
pub fn set_ci(&self, provider: Arc<dyn ci_collector::CiStatsProvider>) {
|
||||
|
|
@ -201,6 +212,9 @@ pub fn start_dashboard(config: DashboardConfig) -> DashboardHandle {
|
|||
let datastore: Arc<Mutex<Option<Arc<dyn datastore_collector::DatastoreStatsProvider>>>> =
|
||||
Arc::new(Mutex::new(None));
|
||||
|
||||
let datastore_factory: Arc<Mutex<Option<Arc<dyn datastore_collector::DatastoreFactory>>>> =
|
||||
Arc::new(Mutex::new(None));
|
||||
|
||||
#[cfg(feature = "ci")]
|
||||
let ci: Arc<Mutex<Option<Arc<dyn ci_collector::CiStatsProvider>>>> =
|
||||
Arc::new(Mutex::new(None));
|
||||
|
|
@ -215,6 +229,7 @@ pub fn start_dashboard(config: DashboardConfig) -> DashboardHandle {
|
|||
#[cfg(feature = "distribution")]
|
||||
Arc::clone(&distribution),
|
||||
Arc::clone(&datastore),
|
||||
Arc::clone(&datastore_factory),
|
||||
#[cfg(feature = "ci")]
|
||||
Arc::clone(&ci),
|
||||
);
|
||||
|
|
@ -260,6 +275,7 @@ pub fn start_dashboard(config: DashboardConfig) -> DashboardHandle {
|
|||
#[cfg(feature = "distribution")]
|
||||
distribution,
|
||||
datastore,
|
||||
datastore_factory,
|
||||
#[cfg(feature = "ci")]
|
||||
ci,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -24,7 +24,7 @@ use crate::distribution_collector::DistributionStatsProvider;
|
|||
#[cfg(feature = "distribution")]
|
||||
use crate::distribution_html::DISTRIBUTION_HTML;
|
||||
|
||||
use crate::datastore_collector::DatastoreStatsProvider;
|
||||
use crate::datastore_collector::{DatastoreFactory, DatastoreStatsProvider, ListScope};
|
||||
use crate::datastore_html::DATASTORE_HTML;
|
||||
|
||||
#[cfg(feature = "ci")]
|
||||
|
|
@ -135,6 +135,7 @@ pub(crate) fn spawn_http_server(
|
|||
#[cfg(feature = "distribution")]
|
||||
distribution: Arc<Mutex<Option<Arc<dyn DistributionStatsProvider>>>>,
|
||||
datastore: Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
datastore_factory: Arc<Mutex<Option<Arc<dyn DatastoreFactory>>>>,
|
||||
#[cfg(feature = "ci")]
|
||||
ci: Arc<Mutex<Option<Arc<dyn CiStatsProvider>>>>,
|
||||
) {
|
||||
|
|
@ -154,6 +155,7 @@ pub(crate) fn spawn_http_server(
|
|||
#[cfg(feature = "distribution")]
|
||||
let distribution = Arc::clone(&distribution);
|
||||
let datastore = Arc::clone(&datastore);
|
||||
let datastore_factory = Arc::clone(&datastore_factory);
|
||||
#[cfg(feature = "ci")]
|
||||
let ci = Arc::clone(&ci);
|
||||
thread::spawn(move || {
|
||||
|
|
@ -165,14 +167,15 @@ pub(crate) fn spawn_http_server(
|
|||
|
||||
let url = request.url().to_string();
|
||||
let path = url.split('?').next().unwrap_or(&url);
|
||||
match path {
|
||||
"/" => respond_html(request, DASHBOARD_HTML, "live"),
|
||||
"/actors" => respond_html(request, ACTORS_HTML, "live"),
|
||||
"/topology" => respond_html(request, TOPOLOGY_HTML, "live"),
|
||||
let method = request.method().as_str();
|
||||
match (method, path) {
|
||||
(_, "/") => respond_html(request, DASHBOARD_HTML, "live"),
|
||||
(_, "/actors") => respond_html(request, ACTORS_HTML, "live"),
|
||||
(_, "/topology") => respond_html(request, TOPOLOGY_HTML, "live"),
|
||||
#[cfg(feature = "distribution")]
|
||||
"/distribution" => respond_html(request, DISTRIBUTION_HTML, "live"),
|
||||
"/datastore" => respond_html(request, DATASTORE_HTML, "live"),
|
||||
"/events" => {
|
||||
(_, "/distribution") => respond_html(request, DISTRIBUTION_HTML, "live"),
|
||||
(_, "/datastore") => respond_html(request, DATASTORE_HTML, "live"),
|
||||
(_, "/events") => {
|
||||
handle_live_sse(
|
||||
request,
|
||||
Arc::clone(&store),
|
||||
|
|
@ -187,24 +190,24 @@ pub(crate) fn spawn_http_server(
|
|||
Arc::clone(&ci),
|
||||
);
|
||||
}
|
||||
"/api/stats" => {
|
||||
(_, "/api/stats") => {
|
||||
handle_stats_api(
|
||||
request,
|
||||
Arc::clone(&runtime),
|
||||
Arc::clone(&collector),
|
||||
);
|
||||
}
|
||||
"/api/history" => {
|
||||
(_, "/api/history") => {
|
||||
handle_history_api(request, Arc::clone(&history));
|
||||
}
|
||||
"/api/topology" => {
|
||||
(_, "/api/topology") => {
|
||||
handle_topology_api(
|
||||
request,
|
||||
Arc::clone(&runtime),
|
||||
Arc::clone(&collector),
|
||||
);
|
||||
}
|
||||
"/api/investigate" => {
|
||||
(_, "/api/investigate") => {
|
||||
handle_investigate_api(
|
||||
request,
|
||||
&url,
|
||||
|
|
@ -214,21 +217,45 @@ pub(crate) fn spawn_http_server(
|
|||
);
|
||||
}
|
||||
#[cfg(feature = "distribution")]
|
||||
"/api/distribution" => {
|
||||
(_, "/api/distribution") => {
|
||||
handle_distribution_api(
|
||||
request,
|
||||
Arc::clone(&distribution),
|
||||
);
|
||||
}
|
||||
"/api/datastore" => {
|
||||
(_, "/api/datastore") => {
|
||||
handle_datastore_api(
|
||||
request,
|
||||
Arc::clone(&datastore),
|
||||
);
|
||||
}
|
||||
"/api/logs" => {
|
||||
(_, "/api/logs") => {
|
||||
handle_logs_api(request, &url, Arc::clone(&store));
|
||||
}
|
||||
("GET", "/api/datastore/list") => {
|
||||
handle_ds_list(request, &url, &datastore);
|
||||
}
|
||||
("GET", "/api/datastore/get") => {
|
||||
handle_ds_get(request, &url, &datastore);
|
||||
}
|
||||
("GET", "/api/datastore/data") => {
|
||||
handle_ds_data(request, &url, &datastore);
|
||||
}
|
||||
("GET", "/api/datastore/status") => {
|
||||
handle_ds_status(request, &datastore);
|
||||
}
|
||||
("POST", "/api/datastore/put") => {
|
||||
handle_ds_put(request, &url, &datastore);
|
||||
}
|
||||
("POST", "/api/datastore/delete") => {
|
||||
handle_ds_delete(request, &url, &datastore);
|
||||
}
|
||||
("POST", "/api/datastore/start") => {
|
||||
handle_ds_start(request, &url, &datastore, &datastore_factory);
|
||||
}
|
||||
("POST", "/api/datastore/shutdown") => {
|
||||
handle_ds_shutdown(request, &datastore);
|
||||
}
|
||||
#[cfg(feature = "ci")]
|
||||
_ if path.starts_with("/api/ci/") => {
|
||||
handle_ci_api(request, path, Arc::clone(&ci));
|
||||
|
|
@ -341,9 +368,21 @@ fn handle_live_sse(
|
|||
// Send datastore snapshot if provider is attached
|
||||
{
|
||||
let maybe_ds = datastore.lock().unwrap().clone();
|
||||
if let Some(provider) = maybe_ds {
|
||||
if let Some(json) = provider.snapshot_json() {
|
||||
if tx.send(format_sse("datastore", &json)).is_err() {
|
||||
match maybe_ds {
|
||||
Some(provider) => {
|
||||
let is_running = provider.is_running();
|
||||
let snap_json = provider.snapshot_json().unwrap_or_else(|| "null".into());
|
||||
let envelope = format!(
|
||||
r#"{{"is_running":{},"snapshot":{}}}"#,
|
||||
is_running, snap_json
|
||||
);
|
||||
if tx.send(format_sse("datastore", &envelope)).is_err() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
None => {
|
||||
let envelope = r#"{"is_running":false,"snapshot":null}"#;
|
||||
if tx.send(format_sse("datastore", envelope)).is_err() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
|
@ -601,6 +640,215 @@ fn parse_query_string(url: &str) -> HashMap<String, String> {
|
|||
params
|
||||
}
|
||||
|
||||
// ── Datastore CRUD API handlers ──────────────────────────────────────────
|
||||
|
||||
fn ds_respond_json(request: tiny_http::Request, json: &str) {
|
||||
let response = tiny_http::Response::from_string(json).with_header(
|
||||
"Content-Type: application/json"
|
||||
.parse::<tiny_http::Header>()
|
||||
.unwrap(),
|
||||
);
|
||||
let _ = request.respond(response);
|
||||
}
|
||||
|
||||
fn ds_respond_bytes(request: tiny_http::Request, data: &[u8]) {
|
||||
let response = tiny_http::Response::from_data(data.to_vec()).with_header(
|
||||
"Content-Type: application/octet-stream"
|
||||
.parse::<tiny_http::Header>()
|
||||
.unwrap(),
|
||||
);
|
||||
let _ = request.respond(response);
|
||||
}
|
||||
|
||||
fn ds_respond_error(request: tiny_http::Request, status: u16, msg: &str) {
|
||||
let json = serde_json::json!({ "error": msg }).to_string();
|
||||
let response = tiny_http::Response::from_string(json)
|
||||
.with_status_code(status)
|
||||
.with_header(
|
||||
"Content-Type: application/json"
|
||||
.parse::<tiny_http::Header>()
|
||||
.unwrap(),
|
||||
);
|
||||
let _ = request.respond(response);
|
||||
}
|
||||
|
||||
fn get_ds_provider(
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
) -> Option<Arc<dyn DatastoreStatsProvider>> {
|
||||
datastore.lock().unwrap().clone()
|
||||
}
|
||||
|
||||
fn handle_ds_list(
|
||||
request: tiny_http::Request,
|
||||
url: &str,
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
) {
|
||||
let provider = match get_ds_provider(datastore) {
|
||||
Some(p) => p,
|
||||
None => { ds_respond_error(request, 503, "datastore not running"); return; }
|
||||
};
|
||||
let params = parse_query_string(url);
|
||||
let scope = match params.get("scope").map(|s| s.as_str()) {
|
||||
Some("local") => ListScope::Local,
|
||||
_ => ListScope::Swarm,
|
||||
};
|
||||
let name_filter = params.get("name").map(|s| s.as_str());
|
||||
match provider.list_objects(name_filter, scope) {
|
||||
Ok(json) => ds_respond_json(request, &json),
|
||||
Err(e) => ds_respond_error(request, 500, &e),
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_ds_get(
|
||||
request: tiny_http::Request,
|
||||
url: &str,
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
) {
|
||||
let provider = match get_ds_provider(datastore) {
|
||||
Some(p) => p,
|
||||
None => { ds_respond_error(request, 503, "datastore not running"); return; }
|
||||
};
|
||||
let params = parse_query_string(url);
|
||||
let hash = match params.get("hash") {
|
||||
Some(h) => h.as_str(),
|
||||
None => { ds_respond_error(request, 400, "missing ?hash= parameter"); return; }
|
||||
};
|
||||
match provider.get_object(hash) {
|
||||
Ok(json) => ds_respond_json(request, &json),
|
||||
Err(e) if e.contains("not found") => ds_respond_error(request, 404, &e),
|
||||
Err(e) => ds_respond_error(request, 500, &e),
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_ds_data(
|
||||
request: tiny_http::Request,
|
||||
url: &str,
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
) {
|
||||
let provider = match get_ds_provider(datastore) {
|
||||
Some(p) => p,
|
||||
None => { ds_respond_error(request, 503, "datastore not running"); return; }
|
||||
};
|
||||
let params = parse_query_string(url);
|
||||
let hash = match params.get("hash") {
|
||||
Some(h) => h.as_str(),
|
||||
None => { ds_respond_error(request, 400, "missing ?hash= parameter"); return; }
|
||||
};
|
||||
match provider.get_data(hash) {
|
||||
Ok(data) => ds_respond_bytes(request, &data),
|
||||
Err(e) if e.contains("not found") => ds_respond_error(request, 404, &e),
|
||||
Err(e) => ds_respond_error(request, 500, &e),
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_ds_status(
|
||||
request: tiny_http::Request,
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
) {
|
||||
let provider = match get_ds_provider(datastore) {
|
||||
Some(p) => p,
|
||||
None => { ds_respond_error(request, 503, "datastore not running"); return; }
|
||||
};
|
||||
match provider.node_status() {
|
||||
Ok(json) => ds_respond_json(request, &json),
|
||||
Err(e) => ds_respond_error(request, 500, &e),
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_ds_put(
|
||||
mut request: tiny_http::Request,
|
||||
url: &str,
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
) {
|
||||
let provider = match get_ds_provider(datastore) {
|
||||
Some(p) => p,
|
||||
None => { ds_respond_error(request, 503, "datastore not running"); return; }
|
||||
};
|
||||
let params = parse_query_string(url);
|
||||
let name = params.get("name").cloned();
|
||||
|
||||
let mut body = Vec::new();
|
||||
if request.as_reader().read_to_end(&mut body).is_err() {
|
||||
return;
|
||||
}
|
||||
|
||||
match provider.put_data(body, name) {
|
||||
Ok(json) => ds_respond_json(request, &json),
|
||||
Err(e) => ds_respond_error(request, 500, &e),
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_ds_delete(
|
||||
request: tiny_http::Request,
|
||||
url: &str,
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
) {
|
||||
let provider = match get_ds_provider(datastore) {
|
||||
Some(p) => p,
|
||||
None => { ds_respond_error(request, 503, "datastore not running"); return; }
|
||||
};
|
||||
let params = parse_query_string(url);
|
||||
let hash = match params.get("hash") {
|
||||
Some(h) => h.as_str(),
|
||||
None => { ds_respond_error(request, 400, "missing ?hash= parameter"); return; }
|
||||
};
|
||||
match provider.delete_object(hash) {
|
||||
Ok(json) => ds_respond_json(request, &json),
|
||||
Err(e) if e.contains("not found") => ds_respond_error(request, 404, &e),
|
||||
Err(e) => ds_respond_error(request, 500, &e),
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_ds_start(
|
||||
request: tiny_http::Request,
|
||||
url: &str,
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
factory: &Arc<Mutex<Option<Arc<dyn DatastoreFactory>>>>,
|
||||
) {
|
||||
// Check if already running
|
||||
{
|
||||
let ds = datastore.lock().unwrap();
|
||||
if ds.is_some() {
|
||||
ds_respond_error(request, 409, "datastore already running");
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
let fac = match factory.lock().unwrap().clone() {
|
||||
Some(f) => f,
|
||||
None => { ds_respond_error(request, 501, "no datastore factory configured"); return; }
|
||||
};
|
||||
|
||||
let params = parse_query_string(url);
|
||||
let storage_path = params.get("storage_path").cloned();
|
||||
|
||||
match fac.start_datastore(storage_path) {
|
||||
Ok(provider) => {
|
||||
*datastore.lock().unwrap() = Some(provider);
|
||||
ds_respond_json(request, r#"{"ok":true}"#);
|
||||
}
|
||||
Err(e) => ds_respond_error(request, 500, &e),
|
||||
}
|
||||
}
|
||||
|
||||
fn handle_ds_shutdown(
|
||||
request: tiny_http::Request,
|
||||
datastore: &Arc<Mutex<Option<Arc<dyn DatastoreStatsProvider>>>>,
|
||||
) {
|
||||
let provider = match get_ds_provider(datastore) {
|
||||
Some(p) => p,
|
||||
None => { ds_respond_error(request, 503, "datastore not running"); return; }
|
||||
};
|
||||
|
||||
match provider.shutdown_datastore() {
|
||||
Ok(()) => {
|
||||
*datastore.lock().unwrap() = None;
|
||||
ds_respond_json(request, r#"{"ok":true}"#);
|
||||
}
|
||||
Err(e) => ds_respond_error(request, 500, &e),
|
||||
}
|
||||
}
|
||||
|
||||
// ── Replay server ───────────────────────────────────────────────────────
|
||||
|
||||
/// Start a replay HTTP server that serves a pre-recorded trace.
|
||||
|
|
|
|||
22
crates/swactor-node/Cargo.toml
Normal file
22
crates/swactor-node/Cargo.toml
Normal file
|
|
@ -0,0 +1,22 @@
|
|||
[package]
|
||||
name = "swactor-node"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
swactor = { path = "../..", features = ["serde", "tracing", "transport"] }
|
||||
runtime-dashboard = { path = "../runtime-dashboard", features = ["distribution"] }
|
||||
swactor-datastore = { path = "../datastore", features = ["node"] }
|
||||
distribution = { path = "../distribution" }
|
||||
clap = { version = "4", features = ["derive"] }
|
||||
ctrlc = "3"
|
||||
iroh = { version = "0.96", optional = true }
|
||||
|
||||
[features]
|
||||
default = ["iroh"]
|
||||
tcp = ["distribution/tcp"]
|
||||
iroh = ["distribution/iroh", "dep:iroh"]
|
||||
|
||||
[[bin]]
|
||||
name = "swactor-node"
|
||||
path = "src/main.rs"
|
||||
491
crates/swactor-node/src/main.rs
Normal file
491
crates/swactor-node/src/main.rs
Normal file
|
|
@ -0,0 +1,491 @@
|
|||
//! swactor-node — unified distributed node with dashboard and datastore.
|
||||
//!
|
||||
//! Combines distribution, runtime dashboard, and content-addressed datastore
|
||||
//! into a single batteries-included binary. Datastore is on by default
|
||||
//! (in-memory) and can be disabled via `--no-datastore`.
|
||||
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::thread;
|
||||
use std::time::Duration;
|
||||
|
||||
use clap::Parser;
|
||||
|
||||
use swactor::actor::{ActorInterface, Ctx};
|
||||
use swactor::config::RuntimeConfig;
|
||||
use swactor::runtime::Runtime;
|
||||
|
||||
use distribution::node::DistributedNodeConfig;
|
||||
use distribution::snapshot::DistributionNodeSnapshot;
|
||||
use distribution::swim::probe::SwimConfig;
|
||||
|
||||
use runtime_dashboard::collector::StatsCollector;
|
||||
use runtime_dashboard::distribution_collector::DistributionStatsProvider;
|
||||
use runtime_dashboard::{start_dashboard, DashboardConfig};
|
||||
|
||||
use swactor_datastore::actors::{BlobStoreActor, DatastoreNode, MetadataActor};
|
||||
use swactor_datastore::bridge::{DatastoreBridge, DatastoreNodeFactory};
|
||||
use swactor_datastore::messages::MetadataMsg;
|
||||
use swactor_datastore::metrics::DatastoreMetrics;
|
||||
use swactor_datastore::storage::{FilesystemBackend, InMemoryBackend};
|
||||
use swactor_datastore::DatastoreConfig;
|
||||
|
||||
use distribution::types::NodeId;
|
||||
|
||||
// ── CLI ──────────────────────────────────────────────────────────────────
|
||||
|
||||
#[derive(Parser)]
|
||||
#[command(
|
||||
name = "swactor-node",
|
||||
about = "Unified swactor node: distribution + dashboard + datastore"
|
||||
)]
|
||||
struct Args {
|
||||
/// Transport to use: iroh or tcp
|
||||
#[arg(long, default_value = "iroh")]
|
||||
transport: String,
|
||||
|
||||
/// Address to listen on for TCP transport (e.g. 10.0.1.10:7000)
|
||||
#[arg(long)]
|
||||
listen: Option<std::net::SocketAddr>,
|
||||
|
||||
/// Seed node address to join (TCP mode: host:port)
|
||||
#[arg(long)]
|
||||
seed: Option<String>,
|
||||
|
||||
/// Seed node's iroh public key (iroh mode: hex-encoded 32-byte key)
|
||||
#[arg(long)]
|
||||
seed_node_id: Option<String>,
|
||||
|
||||
/// Dashboard HTTP port
|
||||
#[arg(long, default_value = "9090")]
|
||||
dashboard_port: u16,
|
||||
|
||||
/// Number of dummy heartbeat actors to register
|
||||
#[arg(long, default_value = "0")]
|
||||
actors: usize,
|
||||
|
||||
/// Storage directory for persistent datastore (omit for in-memory)
|
||||
#[arg(long)]
|
||||
storage_path: Option<String>,
|
||||
|
||||
/// Disable the datastore entirely
|
||||
#[arg(long)]
|
||||
no_datastore: bool,
|
||||
|
||||
/// Chunk size in bytes
|
||||
#[arg(long, default_value = "1048576")]
|
||||
chunk_size: u32,
|
||||
|
||||
/// GC interval in ticks (each tick is ~100ms)
|
||||
#[arg(long, default_value = "1000")]
|
||||
gc_interval: u64,
|
||||
|
||||
/// Dissemination interval in ticks
|
||||
#[arg(long, default_value = "50")]
|
||||
disseminate_interval: u64,
|
||||
}
|
||||
|
||||
// ── Dummy actor ──────────────────────────────────────────────────────────
|
||||
|
||||
#[derive(Clone)]
|
||||
struct Heartbeat;
|
||||
|
||||
struct HeartbeatActor;
|
||||
|
||||
impl ActorInterface for HeartbeatActor {
|
||||
type Incoming = Heartbeat;
|
||||
type Response = ();
|
||||
|
||||
fn handle(&mut self, _ctx: &Ctx, _msg: Heartbeat) {}
|
||||
}
|
||||
|
||||
// ── Snapshot provider ────────────────────────────────────────────────────
|
||||
|
||||
struct SnapshotProvider {
|
||||
snapshot: Arc<Mutex<Option<DistributionNodeSnapshot>>>,
|
||||
}
|
||||
|
||||
impl DistributionStatsProvider for SnapshotProvider {
|
||||
fn snapshot(&self) -> Option<DistributionNodeSnapshot> {
|
||||
self.snapshot.lock().unwrap().clone()
|
||||
}
|
||||
}
|
||||
|
||||
// ── Main ─────────────────────────────────────────────────────────────────
|
||||
|
||||
fn main() {
|
||||
let args = Args::parse();
|
||||
let stop = Arc::new(AtomicBool::new(false));
|
||||
|
||||
// Signal handler
|
||||
{
|
||||
let stop = Arc::clone(&stop);
|
||||
ctrlc::set_handler(move || {
|
||||
stop.store(true, Ordering::Relaxed);
|
||||
})
|
||||
.expect("failed to set signal handler");
|
||||
}
|
||||
|
||||
// Start dashboard
|
||||
let dash = start_dashboard(DashboardConfig {
|
||||
port: args.dashboard_port,
|
||||
..Default::default()
|
||||
});
|
||||
dash.install_tracing();
|
||||
|
||||
// Create actor runtime
|
||||
let num_threads = 2;
|
||||
let collector = StatsCollector::new(num_threads);
|
||||
let mut rt = Runtime::new(RuntimeConfig {
|
||||
num_threads,
|
||||
max_actors: 1024,
|
||||
channel_buffer_size: 2000,
|
||||
..Default::default()
|
||||
});
|
||||
rt.set_stats_hook(collector.clone());
|
||||
|
||||
let handle = rt.run().expect("failed to start runtime");
|
||||
dash.set_runtime(handle.runtime.clone(), collector);
|
||||
|
||||
// Datastore setup
|
||||
let mut ds_metadata_addr = None;
|
||||
let ds_gc_interval = args.gc_interval;
|
||||
let ds_disseminate_interval = args.disseminate_interval;
|
||||
|
||||
if !args.no_datastore {
|
||||
let node_id = generate_node_id();
|
||||
let config = DatastoreConfig {
|
||||
chunk_size: args.chunk_size,
|
||||
storage_path: args
|
||||
.storage_path
|
||||
.as_ref()
|
||||
.map(|s| s.into())
|
||||
.unwrap_or_else(|| "datastore".into()),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let backend: Box<dyn swactor_datastore::StorageBackend> = match &args.storage_path {
|
||||
Some(path) => {
|
||||
let p = std::path::PathBuf::from(path);
|
||||
std::fs::create_dir_all(&p).expect("failed to create storage directory");
|
||||
Box::new(FilesystemBackend::new(p))
|
||||
}
|
||||
None => Box::new(InMemoryBackend::new()),
|
||||
};
|
||||
|
||||
let blob_store_addr = handle
|
||||
.runtime
|
||||
.spawn(BlobStoreActor::new(backend))
|
||||
.expect("failed to spawn BlobStoreActor");
|
||||
|
||||
let mut metadata = MetadataActor::new(node_id, &config);
|
||||
metadata.set_blob_store(blob_store_addr);
|
||||
let metadata_addr = handle
|
||||
.runtime
|
||||
.spawn(metadata)
|
||||
.expect("failed to spawn MetadataActor");
|
||||
|
||||
let datastore_node =
|
||||
DatastoreNode::new(node_id, blob_store_addr, metadata_addr, config);
|
||||
let datastore_addr = handle
|
||||
.runtime
|
||||
.spawn(datastore_node)
|
||||
.expect("failed to spawn DatastoreNode");
|
||||
|
||||
let node_hex: String = node_id.0.iter().map(|b| format!("{b:02x}")).collect();
|
||||
let metrics = Arc::new(DatastoreMetrics::new());
|
||||
metrics.set_node_id(node_hex.clone());
|
||||
|
||||
let bridge = DatastoreBridge::new(
|
||||
Arc::clone(&metrics),
|
||||
Arc::clone(&handle.runtime),
|
||||
datastore_addr,
|
||||
metadata_addr,
|
||||
blob_store_addr,
|
||||
);
|
||||
|
||||
dash.set_datastore(Arc::new(bridge));
|
||||
ds_metadata_addr = Some(metadata_addr);
|
||||
|
||||
if args.storage_path.is_some() {
|
||||
eprintln!("Datastore: persistent ({})", args.storage_path.as_ref().unwrap());
|
||||
} else {
|
||||
eprintln!("Datastore: in-memory");
|
||||
}
|
||||
} else {
|
||||
eprintln!("Datastore: disabled");
|
||||
}
|
||||
|
||||
// Always wire up the factory so the UI can start/stop datastore
|
||||
let factory = DatastoreNodeFactory::new(
|
||||
Arc::clone(&handle.runtime),
|
||||
args.chunk_size,
|
||||
);
|
||||
dash.set_datastore_factory(Arc::new(factory));
|
||||
|
||||
// Distribution config
|
||||
let swim_config = SwimConfig {
|
||||
probe_interval: 5,
|
||||
probe_timeout: 3,
|
||||
indirect_probes: 2,
|
||||
suspicion_timeout: 20,
|
||||
dead_reprobe_interval: 50,
|
||||
};
|
||||
let node_config = DistributedNodeConfig {
|
||||
swim: swim_config,
|
||||
cache_capacity: 1000,
|
||||
republish_interval: 500,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
match args.transport.as_str() {
|
||||
#[cfg(feature = "iroh")]
|
||||
"iroh" => run_iroh(
|
||||
args,
|
||||
node_config,
|
||||
&handle,
|
||||
&dash,
|
||||
&stop,
|
||||
ds_metadata_addr,
|
||||
ds_gc_interval,
|
||||
ds_disseminate_interval,
|
||||
),
|
||||
#[cfg(feature = "tcp")]
|
||||
"tcp" => run_tcp(
|
||||
args,
|
||||
node_config,
|
||||
&handle,
|
||||
&dash,
|
||||
&stop,
|
||||
ds_metadata_addr,
|
||||
ds_gc_interval,
|
||||
ds_disseminate_interval,
|
||||
),
|
||||
other => {
|
||||
eprintln!("Unknown or unavailable transport: {other}");
|
||||
eprintln!("Available transports:");
|
||||
#[cfg(feature = "iroh")]
|
||||
eprintln!(" iroh");
|
||||
#[cfg(feature = "tcp")]
|
||||
eprintln!(" tcp");
|
||||
std::process::exit(1);
|
||||
}
|
||||
}
|
||||
|
||||
eprintln!("\nShutting down...");
|
||||
handle.shutdown();
|
||||
dash.shutdown();
|
||||
handle.join();
|
||||
}
|
||||
|
||||
// ── TCP transport ────────────────────────────────────────────────────────
|
||||
|
||||
#[cfg(feature = "tcp")]
|
||||
fn run_tcp(
|
||||
args: Args,
|
||||
node_config: DistributedNodeConfig,
|
||||
handle: &swactor::runtime::RuntimeHandle,
|
||||
dash: &runtime_dashboard::DashboardHandle,
|
||||
stop: &Arc<AtomicBool>,
|
||||
ds_metadata_addr: Option<swactor::actor::ActorAddress>,
|
||||
ds_gc_interval: u64,
|
||||
ds_disseminate_interval: u64,
|
||||
) {
|
||||
use distribution::driver::NodeDriver;
|
||||
|
||||
let listen_addr = args.listen.expect("--listen is required for TCP mode");
|
||||
let mut driver =
|
||||
NodeDriver::new(listen_addr, node_config).expect("failed to create node driver");
|
||||
|
||||
eprintln!(
|
||||
"Node {} listening on {} (TCP)",
|
||||
hex(&driver.node_id().0[..4]),
|
||||
driver.listen_addr(),
|
||||
);
|
||||
|
||||
// Join seed if provided
|
||||
if let Some(seed) = args.seed {
|
||||
let seed_addr: std::net::SocketAddr = seed.parse().expect("invalid seed address");
|
||||
eprintln!("Joining cluster via seed {seed_addr}");
|
||||
driver.join(&[seed_addr]);
|
||||
}
|
||||
|
||||
// Spawn and register actors
|
||||
let actor_addrs = spawn_actors(args.actors, handle, driver.node_mut());
|
||||
|
||||
// Wire distribution snapshot to dashboard
|
||||
let cached_snapshot: Arc<Mutex<Option<DistributionNodeSnapshot>>> =
|
||||
Arc::new(Mutex::new(Some(driver.snapshot())));
|
||||
let provider = SnapshotProvider {
|
||||
snapshot: Arc::clone(&cached_snapshot),
|
||||
};
|
||||
dash.set_distribution(Arc::new(provider));
|
||||
|
||||
eprintln!("Dashboard at http://0.0.0.0:{}", args.dashboard_port);
|
||||
|
||||
// Main loop
|
||||
let mut round: u64 = 0;
|
||||
while !stop.load(Ordering::Relaxed) {
|
||||
round += 1;
|
||||
|
||||
driver.recv();
|
||||
driver.tick();
|
||||
|
||||
for addr in &actor_addrs {
|
||||
let _ = handle.runtime.send_to(*addr, Heartbeat);
|
||||
}
|
||||
|
||||
*cached_snapshot.lock().unwrap() = Some(driver.snapshot());
|
||||
|
||||
// Datastore ticks
|
||||
if let Some(metadata_addr) = ds_metadata_addr {
|
||||
if round % ds_gc_interval == 0 {
|
||||
let _ = handle.runtime.send_to(metadata_addr, MetadataMsg::GcTick);
|
||||
}
|
||||
if round % ds_disseminate_interval == 0 {
|
||||
let _ = handle
|
||||
.runtime
|
||||
.send_to(metadata_addr, MetadataMsg::DisseminateTick);
|
||||
}
|
||||
}
|
||||
|
||||
thread::sleep(Duration::from_millis(100));
|
||||
}
|
||||
}
|
||||
|
||||
// ── iroh transport ───────────────────────────────────────────────────────
|
||||
|
||||
#[cfg(feature = "iroh")]
|
||||
fn run_iroh(
|
||||
args: Args,
|
||||
node_config: DistributedNodeConfig,
|
||||
handle: &swactor::runtime::RuntimeHandle,
|
||||
dash: &runtime_dashboard::DashboardHandle,
|
||||
stop: &Arc<AtomicBool>,
|
||||
ds_metadata_addr: Option<swactor::actor::ActorAddress>,
|
||||
ds_gc_interval: u64,
|
||||
ds_disseminate_interval: u64,
|
||||
) {
|
||||
use distribution::iroh_driver::{IrohDriver, IrohDriverConfig};
|
||||
use iroh::RelayMode;
|
||||
|
||||
let iroh_config = IrohDriverConfig {
|
||||
secret_key: None,
|
||||
relay_mode: RelayMode::Default,
|
||||
node: node_config,
|
||||
};
|
||||
let mut driver = IrohDriver::new(iroh_config).expect("failed to create iroh driver");
|
||||
|
||||
eprintln!("Node {} started (iroh)", hex(&driver.node_id().0[..4]));
|
||||
|
||||
// Join seed if provided
|
||||
if let Some(seed_hex) = args.seed_node_id {
|
||||
let seed_bytes = hex_to_bytes(&seed_hex).expect("invalid seed node ID hex");
|
||||
let seed_key =
|
||||
iroh::PublicKey::from_bytes(&seed_bytes).expect("invalid seed public key");
|
||||
eprintln!("Joining cluster via seed {}", &seed_hex[..8]);
|
||||
driver.join(&[seed_key]);
|
||||
}
|
||||
|
||||
// Spawn and register actors
|
||||
let actor_addrs = spawn_actors(args.actors, handle, driver.node_mut());
|
||||
|
||||
// Wire distribution snapshot to dashboard
|
||||
let cached_snapshot: Arc<Mutex<Option<DistributionNodeSnapshot>>> =
|
||||
Arc::new(Mutex::new(Some(driver.snapshot())));
|
||||
let provider = SnapshotProvider {
|
||||
snapshot: Arc::clone(&cached_snapshot),
|
||||
};
|
||||
dash.set_distribution(Arc::new(provider));
|
||||
|
||||
eprintln!("Dashboard at http://0.0.0.0:{}", args.dashboard_port);
|
||||
|
||||
// Main loop
|
||||
let mut round: u64 = 0;
|
||||
while !stop.load(Ordering::Relaxed) {
|
||||
round += 1;
|
||||
|
||||
driver.recv();
|
||||
driver.tick();
|
||||
|
||||
for addr in &actor_addrs {
|
||||
let _ = handle.runtime.send_to(*addr, Heartbeat);
|
||||
}
|
||||
|
||||
*cached_snapshot.lock().unwrap() = Some(driver.snapshot());
|
||||
|
||||
// Datastore ticks
|
||||
if let Some(metadata_addr) = ds_metadata_addr {
|
||||
if round % ds_gc_interval == 0 {
|
||||
let _ = handle.runtime.send_to(metadata_addr, MetadataMsg::GcTick);
|
||||
}
|
||||
if round % ds_disseminate_interval == 0 {
|
||||
let _ = handle
|
||||
.runtime
|
||||
.send_to(metadata_addr, MetadataMsg::DisseminateTick);
|
||||
}
|
||||
}
|
||||
|
||||
thread::sleep(Duration::from_millis(100));
|
||||
}
|
||||
|
||||
driver.shutdown();
|
||||
}
|
||||
|
||||
// ── Helpers ──────────────────────────────────────────────────────────────
|
||||
|
||||
fn spawn_actors(
|
||||
count: usize,
|
||||
handle: &swactor::runtime::RuntimeHandle,
|
||||
node: &mut distribution::node::DistributedNode,
|
||||
) -> Vec<swactor::actor::ActorAddress> {
|
||||
let mut addrs = Vec::new();
|
||||
for _ in 0..count {
|
||||
match handle.runtime.spawn(HeartbeatActor) {
|
||||
Ok(addr) => {
|
||||
node.register_actor(addr, 1);
|
||||
addrs.push(addr);
|
||||
}
|
||||
Err(e) => eprintln!("failed to spawn actor: {e}"),
|
||||
}
|
||||
}
|
||||
if !addrs.is_empty() {
|
||||
eprintln!("Registered {} actors", addrs.len());
|
||||
}
|
||||
addrs
|
||||
}
|
||||
|
||||
fn generate_node_id() -> NodeId {
|
||||
let mut bytes = [0u8; 32];
|
||||
for (i, b) in std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_nanos()
|
||||
.to_le_bytes()
|
||||
.iter()
|
||||
.enumerate()
|
||||
{
|
||||
bytes[i % 32] ^= *b;
|
||||
}
|
||||
let pid = std::process::id();
|
||||
for (i, b) in pid.to_le_bytes().iter().enumerate() {
|
||||
bytes[i + 16] ^= *b;
|
||||
}
|
||||
NodeId(bytes)
|
||||
}
|
||||
|
||||
fn hex(bytes: &[u8]) -> String {
|
||||
bytes.iter().map(|b| format!("{b:02x}")).collect()
|
||||
}
|
||||
|
||||
#[cfg(feature = "iroh")]
|
||||
fn hex_to_bytes(hex: &str) -> Option<[u8; 32]> {
|
||||
if hex.len() != 64 {
|
||||
return None;
|
||||
}
|
||||
let mut bytes = [0u8; 32];
|
||||
for i in 0..32 {
|
||||
bytes[i] = u8::from_str_radix(&hex[i * 2..i * 2 + 2], 16).ok()?;
|
||||
}
|
||||
Some(bytes)
|
||||
}
|
||||
|
|
@ -66,6 +66,14 @@ enum Cmd {
|
|||
/// Build the crypto WASM module
|
||||
Wasm,
|
||||
|
||||
/// Launch a dev node (distribution + dashboard + datastore)
|
||||
#[command(trailing_var_arg = true)]
|
||||
DevNode {
|
||||
/// Extra arguments forwarded to the dev node
|
||||
#[arg(allow_hyphen_values = true)]
|
||||
extra: Vec<String>,
|
||||
},
|
||||
|
||||
/// Run a datastore CLI command
|
||||
#[command(trailing_var_arg = true)]
|
||||
Cli {
|
||||
|
|
@ -236,9 +244,13 @@ fn run_step(group_name: &str, step: &TestStep) -> bool {
|
|||
fn print_usage() {
|
||||
println!(
|
||||
"\
|
||||
USAGE: cargo xtask test <GROUP>
|
||||
USAGE: cargo xtask <COMMAND>
|
||||
|
||||
GROUPS:
|
||||
COMMANDS:
|
||||
test <GROUP> Run a test group
|
||||
dev-node [OPTS] Launch a dev node (distribution + dashboard + datastore)
|
||||
|
||||
TEST GROUPS:
|
||||
core Actor runtime, message delivery, property tests
|
||||
distribution Distribution protocol + datastore
|
||||
cluster-sims Deterministic cluster simulations
|
||||
|
|
@ -246,8 +258,17 @@ GROUPS:
|
|||
essential core + distribution + integrated (merge gate)
|
||||
all Every test group
|
||||
|
||||
FLAGS:
|
||||
--list Show all groups and the cargo commands they run"
|
||||
TEST FLAGS:
|
||||
--list Show all groups and the cargo commands they run
|
||||
|
||||
DEV OPTIONS:
|
||||
--port PORT Dashboard port (default: 9090)
|
||||
--actors N Dummy heartbeat actors (default: 3)
|
||||
--storage PATH Persistent storage dir (omit for in-memory)
|
||||
--no-datastore Disable datastore entirely
|
||||
--tcp Use TCP transport instead of iroh (requires --listen)
|
||||
--listen ADDR TCP listen address (e.g. 127.0.0.1:7000)
|
||||
--release Build in release mode"
|
||||
);
|
||||
}
|
||||
|
||||
|
|
@ -325,6 +346,128 @@ fn run_test(group: Option<String>, list: bool) {
|
|||
);
|
||||
}
|
||||
|
||||
// ── Dev node launcher ───────────────────────────────────────────────────
|
||||
|
||||
fn run_dev(extra_args: Vec<String>) {
|
||||
ignore_sigint();
|
||||
let mut port = "9090".to_string();
|
||||
let mut listen: Option<String> = None;
|
||||
let mut actors = "3".to_string();
|
||||
let mut storage: Option<String> = None;
|
||||
let mut no_datastore = false;
|
||||
let mut use_tcp = false;
|
||||
let mut release = false;
|
||||
|
||||
let mut i = 0;
|
||||
while i < extra_args.len() {
|
||||
match extra_args[i].as_str() {
|
||||
"--port" => {
|
||||
i += 1;
|
||||
port = extra_args.get(i).cloned().unwrap_or_else(|| {
|
||||
eprintln!("--port requires a value");
|
||||
std::process::exit(1);
|
||||
});
|
||||
}
|
||||
"--listen" => {
|
||||
i += 1;
|
||||
listen = Some(extra_args.get(i).cloned().unwrap_or_else(|| {
|
||||
eprintln!("--listen requires a value");
|
||||
std::process::exit(1);
|
||||
}));
|
||||
}
|
||||
"--actors" => {
|
||||
i += 1;
|
||||
actors = extra_args.get(i).cloned().unwrap_or_else(|| {
|
||||
eprintln!("--actors requires a value");
|
||||
std::process::exit(1);
|
||||
});
|
||||
}
|
||||
"--storage" => {
|
||||
i += 1;
|
||||
storage = Some(extra_args.get(i).cloned().unwrap_or_else(|| {
|
||||
eprintln!("--storage requires a value");
|
||||
std::process::exit(1);
|
||||
}));
|
||||
}
|
||||
"--no-datastore" => {
|
||||
no_datastore = true;
|
||||
}
|
||||
"--tcp" => {
|
||||
use_tcp = true;
|
||||
}
|
||||
"--release" => {
|
||||
release = true;
|
||||
}
|
||||
other => {
|
||||
eprintln!("Unknown dev-node option: {other}");
|
||||
std::process::exit(1);
|
||||
}
|
||||
}
|
||||
i += 1;
|
||||
}
|
||||
|
||||
if use_tcp && listen.is_none() {
|
||||
listen = Some("127.0.0.1:7000".to_string());
|
||||
}
|
||||
|
||||
let mut cargo_args: Vec<&str> = vec!["run", "-p", "swactor-node"];
|
||||
if use_tcp {
|
||||
cargo_args.push("--features");
|
||||
cargo_args.push("tcp");
|
||||
}
|
||||
if release {
|
||||
cargo_args.push("--release");
|
||||
}
|
||||
cargo_args.push("--");
|
||||
|
||||
if use_tcp {
|
||||
cargo_args.push("--transport");
|
||||
cargo_args.push("tcp");
|
||||
}
|
||||
|
||||
let listen_ref;
|
||||
if let Some(ref l) = listen {
|
||||
listen_ref = l.as_str();
|
||||
cargo_args.push("--listen");
|
||||
cargo_args.push(listen_ref);
|
||||
}
|
||||
|
||||
cargo_args.push("--dashboard-port");
|
||||
cargo_args.push(&port);
|
||||
cargo_args.push("--actors");
|
||||
cargo_args.push(&actors);
|
||||
|
||||
let storage_ref;
|
||||
if let Some(ref s) = storage {
|
||||
storage_ref = s.as_str();
|
||||
cargo_args.push("--storage-path");
|
||||
cargo_args.push(storage_ref);
|
||||
}
|
||||
|
||||
if no_datastore {
|
||||
cargo_args.push("--no-datastore");
|
||||
}
|
||||
|
||||
println!(" cargo {}", cargo_args.join(" "));
|
||||
println!();
|
||||
|
||||
let status = Command::new("cargo")
|
||||
.args(&cargo_args)
|
||||
.status();
|
||||
|
||||
match status {
|
||||
Ok(s) => {
|
||||
if !s.success() {
|
||||
std::process::exit(s.code().unwrap_or(1));
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
eprintln!("Failed to execute cargo: {e}");
|
||||
std::process::exit(1);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn run_node(
|
||||
port: Option<u16>,
|
||||
storage_path: Option<String>,
|
||||
|
|
@ -515,6 +658,7 @@ fn main() {
|
|||
auth_dir,
|
||||
extra,
|
||||
} => run_node(port, storage_path, auth, auth_dir, extra, &config.node),
|
||||
Cmd::DevNode { extra } => run_dev(extra),
|
||||
Cmd::Cli {
|
||||
url,
|
||||
key,
|
||||
|
|
|
|||
Loading…
Reference in a new issue