use std::collections::HashMap;
use std::sync::Arc;
use axum::body::Bytes;
use axum::extract::{Query, State};
use axum::http::{HeaderValue, StatusCode};
use axum::response::Response;
use axum::Extension;
use recall_wire::discovery::{self, Auth, Build, Protocol, ServerInfo};
use recall_wire::{content_sha256, Discovery, PROTOCOL};
use recall_wire::{
AdminStats, ClaudeCliStatus, Health, MergeError, MergeSide, MergeStatus, PushRequest,
PushResponse, SyncResponse, WorkerStatus,
};
use super::auth::{Caller, SignedRequestInfo};
use super::respond::{error, internal, json};
use super::AppState;
use crate::audit::leaf;
use crate::now;
use crate::store::Queued;
const REQUIRED_FIELDS_MSG: &str =
"project_key, file_path, and content (string) are required, unless deleted is true";
pub(super) const WORKER_STALE: std::time::Duration = std::time::Duration::from_secs(120);
pub(super) async fn handle_push(
State(state): State<Arc<AppState>>,
caller: Option<Extension<Caller>>,
signed: Option<Extension<SignedRequestInfo>>,
body: Bytes,
) -> Response {
let mut req: PushRequest = match serde_json::from_slice(&body) {
Ok(req) => req,
Err(_) => {
let missing_field = serde_json::from_slice::<serde_json::Value>(&body)
.ok()
.and_then(|v| {
v.as_object()
.map(|o| !o.contains_key("project_key") || !o.contains_key("file_path"))
})
.unwrap_or(false);
return error(
StatusCode::BAD_REQUEST,
if missing_field {
REQUIRED_FIELDS_MSG
} else {
"invalid json body"
},
);
}
};
if req.project_key.is_empty() || req.file_path.is_empty() {
return error(StatusCode::BAD_REQUEST, REQUIRED_FIELDS_MSG);
}
match recall_wire::validate_file_path(&req.file_path) {
Ok(()) => {}
Err(e @ recall_wire::ValidationError::FilePathTooLong) => {
return error(StatusCode::BAD_REQUEST, &e.to_string())
}
Err(_) => {
return error(
StatusCode::BAD_REQUEST,
"file_path must be relative, no traversal",
)
}
}
if let Err(e) = recall_wire::validate_project_key(&req.project_key) {
return error(StatusCode::BAD_REQUEST, &e.to_string());
}
if let Some(base) = &mut req.base_sha256 {
if let Err(e) = recall_wire::validate_base_sha256(base) {
return error(StatusCode::BAD_REQUEST, &e.to_string());
}
base.make_ascii_lowercase();
}
if let Some(Extension(Caller::Device { name, .. })) = &caller {
req.source_env.clone_from(name);
}
let caller_ref = caller.as_ref().map(|Extension(c)| c);
let signed_ref = signed.as_ref().map(|Extension(s)| s);
let actor = super::audit::actor_for(caller_ref);
let request = super::audit::signed_request_for(signed_ref, None);
if req.deleted {
let result = state.store.tombstone_audited(
&req.project_key,
&req.file_path,
&req.source_env,
|seq, at| {
leaf::encode(
seq,
at,
leaf::action::DELETE,
&actor,
leaf::subject_file(&leaf::FileChange {
project_key: &req.project_key,
file_path: &req.file_path,
deleted: true,
stored_sha256: &content_sha256(""),
base_sha256: req.base_sha256.as_deref(),
merged: false,
merge_job: None,
}),
request.as_ref(),
)
},
);
let updated_at = match result {
Ok(at) => at,
Err(e) => return internal(e),
};
return json(
StatusCode::OK,
&PushResponse {
ok: true,
project_key: req.project_key,
file_path: req.file_path,
deleted: true,
merged: false,
updated_at,
merge_job: None,
},
);
}
let existing = match state.store.get(&req.project_key, &req.file_path) {
Ok(e) => e,
Err(e) => return internal(e),
};
let Some(incoming) = req.content.clone() else {
return error(StatusCode::BAD_REQUEST, REQUIRED_FIELDS_MSG);
};
let mut content = incoming.clone();
let mut merged = false;
let stale = needs_merge(
&state,
existing.as_ref(),
&incoming,
req.base_sha256.as_deref(),
);
if stale.is_some() {
match state.store.enrolled_worker() {
Ok(Some(_)) => match merge_here_instead(&state) {
Ok(None) => return queue_merge(&state, req, incoming, &actor, request.as_ref()),
Ok(Some(why)) => eprintln!(
"{why}, so {}/{} is merged here rather than queued",
req.project_key, req.file_path
),
Err(e) => return internal(e),
},
Ok(None) => {}
Err(e) => return internal(e),
}
}
if let Some(stored) = stale.filter(|_| state.read().claude_status.logged_in) {
match state.merger.merge(&stored.content, &incoming).await {
Ok(out) => {
content = out;
merged = true;
let mut rt = state.write();
rt.last_merge_at = now();
rt.last_merge_error = None;
}
Err(e) => {
eprintln!(
"merge failed for {}/{}, falling back to last-write-wins: {e}",
req.project_key, req.file_path
);
state.write().last_merge_error = Some(MergeError {
message: e.to_string(),
at: now(),
});
}
}
}
let result = state.store.upsert_audited(
&req.project_key,
&req.file_path,
&content,
&req.source_env,
|seq, at| {
leaf::encode(
seq,
at,
leaf::action::PUSH,
&actor,
leaf::subject_file(&leaf::FileChange {
project_key: &req.project_key,
file_path: &req.file_path,
deleted: false,
stored_sha256: &content_sha256(&content),
base_sha256: req.base_sha256.as_deref(),
merged,
merge_job: None,
}),
request.as_ref(),
)
},
);
let updated_at = match result {
Ok(at) => at,
Err(e) => return internal(e),
};
json(
StatusCode::OK,
&PushResponse {
ok: true,
project_key: req.project_key,
file_path: req.file_path,
deleted: false,
merged,
updated_at,
merge_job: None,
},
)
}
fn queue_merge(
state: &AppState,
req: PushRequest,
incoming: String,
actor: &leaf::Actor<'_>,
request: Option<&leaf::SignedRequest<'_>>,
) -> Response {
let job_id = match super::devices::new_id("job_", 10) {
Ok(id) => id,
Err(e) => return internal(e),
};
let side = MergeSide {
sha256: recall_wire::content_sha256(&incoming),
content: incoming,
source_env: req.source_env.clone(),
updated_at: String::new(),
};
let queued = state.store.write_and_queue_merge_audited(
&req.project_key,
&req.file_path,
&side,
&job_id,
time::OffsetDateTime::now_utc(),
|seq, at, queued| {
leaf::encode(
seq,
at,
leaf::action::PUSH,
actor,
leaf::subject_file(&leaf::FileChange {
project_key: &req.project_key,
file_path: &req.file_path,
deleted: false,
stored_sha256: &side.sha256,
base_sha256: req.base_sha256.as_deref(),
merged: false,
merge_job: match queued {
Queued::Queued(id) => Some(id),
Queued::Nothing | Queued::Full => None,
},
}),
request,
)
},
);
let (queued, updated_at) = match queued {
Ok(done) => done,
Err(e) => return internal(e),
};
let merge_job = match queued {
Queued::Queued(id) => {
state.jobs_ready.notify_waiters();
Some(id)
}
Queued::Nothing => None,
Queued::Full => {
let message = format!(
"the merge queue is full ({} jobs), so a conflicting push was stored \
last-write-wins; is the worker running?",
crate::store::MAX_OPEN_JOBS,
);
eprintln!("{message} ({}/{})", req.project_key, req.file_path);
state.write().last_merge_error = Some(MergeError { message, at: now() });
None
}
};
json(
StatusCode::OK,
&PushResponse {
ok: true,
project_key: req.project_key,
file_path: req.file_path,
deleted: false,
merged: false,
updated_at,
merge_job,
},
)
}
fn merge_here_instead(state: &AppState) -> anyhow::Result<Option<String>> {
if !state.read().claude_status.logged_in {
return Ok(None);
}
super::jobs::expire_leases(state)?;
let queue = state.store.queue_status()?;
if (queue.queued + queue.leased) as usize >= crate::store::MAX_OPEN_JOBS {
return Ok(Some(format!(
"the merge queue is full ({} jobs)",
crate::store::MAX_OPEN_JOBS
)));
}
let quiet = state.read().worker_last_claim.elapsed();
if quiet >= WORKER_STALE && queue.leased == 0 {
return Ok(Some(format!(
"the merge worker has not asked for work in {}s",
quiet.as_secs()
)));
}
Ok(None)
}
fn needs_merge<'a>(
state: &AppState,
existing: Option<&'a crate::store::Existing>,
incoming: &str,
base_sha256: Option<&str>,
) -> Option<&'a crate::store::Existing> {
if !state.cfg.merge_enabled {
return None;
}
let stored = existing.filter(|e| !e.deleted && e.content != incoming)?;
if base_sha256.is_some_and(|base| {
base.eq_ignore_ascii_case(&recall_wire::content_sha256(&stored.content))
}) {
return None;
}
Some(stored)
}
pub(super) async fn handle_pull(
State(state): State<Arc<AppState>>,
caller: Option<Extension<Caller>>,
signed: Option<Extension<SignedRequestInfo>>,
Query(params): Query<HashMap<String, String>>,
) -> Response {
let Some(project_key) = params.get("project_key").filter(|k| !k.is_empty()) else {
return error(
StatusCode::BAD_REQUEST,
"project_key query param is required",
);
};
if let Err(e) = recall_wire::validate_project_key(project_key) {
return error(StatusCode::BAD_REQUEST, &e.to_string());
}
let files = match state.store.list(project_key) {
Ok(files) => files,
Err(e) => return internal(e),
};
let caller_ref = caller.as_ref().map(|Extension(c)| c);
let signed_ref = signed.as_ref().map(|Extension(s)| s);
let actor = super::audit::actor_for(caller_ref);
let request = super::audit::signed_request_for(signed_ref, None);
if let Err(e) = state.store.audit_append(|seq, at| {
leaf::encode(
seq,
at,
leaf::action::PULL,
&actor,
leaf::subject_pull(project_key),
request.as_ref(),
)
}) {
return internal(e);
}
let mut resp = json(
StatusCode::OK,
&SyncResponse {
project_key: project_key.clone(),
files,
},
);
let (tree_size, root) = state.store.audit_checkpoint();
let checkpoint = recall_wire::AuditCheckpoint {
tree_size,
root_hash: super::audit::base64_hash(&root),
};
if let Ok(value) = HeaderValue::from_str(&checkpoint.to_header_value()) {
resp.headers_mut()
.insert(recall_wire::audit::CHECKPOINT_HEADER, value);
}
resp
}
fn last_offbox_at(backup_dir: &str) -> String {
if backup_dir.is_empty() {
return String::new();
}
std::fs::read_to_string(std::path::Path::new(backup_dir).join(".last-offbox"))
.map(|s| s.trim().to_string())
.unwrap_or_default()
}
pub(super) async fn handle_health(State(state): State<Arc<AppState>>) -> Response {
let last_sync_at = match state.store.last_sync_at() {
Ok(v) => v,
Err(e) => return internal(e),
};
let worker = match state.store.enrolled_worker() {
Ok(w) => w,
Err(e) => return internal(e),
};
if let Err(e) = super::jobs::expire_leases(&state) {
return internal(e);
}
let queue = match state.store.queue_status() {
Ok(q) if worker.is_some() || q.queued + q.leased + q.failed > 0 => Some(q),
Ok(_) => None,
Err(e) => return internal(e),
};
let rt = state.read();
let claude_cli = if worker.is_some() {
rt.worker_cli.clone().unwrap_or_default()
} else if rt.claude_status.checked_at.is_empty() {
ClaudeCliStatus::default()
} else {
ClaudeCliStatus {
checked_at: rt.claude_status.checked_at.clone(),
available: Some(rt.claude_status.available),
logged_in: Some(rt.claude_status.logged_in),
error: rt.claude_status.error.clone(),
}
};
let worker = worker.map(|device| WorkerStatus {
last_claim_at: rt.worker_last_claim_at.clone(),
agent: if rt.worker_agent.is_empty() {
device.agent
} else {
rt.worker_agent.clone()
},
});
let body = Health {
status: "ok".to_string(),
git_commit: state.cfg.git_commit.clone(),
started_at: state.started_at.clone(),
last_sync_at,
last_backup_at: rt.last_backup_at.clone(),
last_offbox_at: last_offbox_at(&state.cfg.backup_dir),
merge: MergeStatus {
enabled: state.cfg.merge_enabled,
claude_cli,
last_merge_at: rt.last_merge_at.clone(),
last_merge_error: rt.last_merge_error.clone(),
worker,
queue,
},
};
drop(rt);
json(StatusCode::OK, &body)
}
const MIN_CLIENT: &str = "0.1.0";
pub(super) async fn handle_discovery(State(state): State<Arc<AppState>>) -> Response {
let revision = Some(state.cfg.git_commit.clone())
.filter(|c| !c.is_empty() && c != "unknown")
.or_else(|| discovery::revision().map(str::to_string));
let channel = discovery::channel();
let mut capabilities = std::collections::BTreeMap::new();
capabilities.insert(
discovery::CAPABILITY_DEVICES.to_string(),
serde_json::to_value(super::devices::capability()).unwrap_or_default(),
);
capabilities.insert("merge_base".to_string(), serde_json::json!({}));
capabilities.insert(
discovery::CAPABILITY_MERGE_QUEUE.to_string(),
serde_json::json!({}),
);
capabilities.insert(
discovery::CAPABILITY_AUDIT.to_string(),
serde_json::to_value(recall_wire::AuditCapability {
leaf_version: recall_wire::audit::LEAF_VERSION,
max_page: recall_wire::audit::MAX_PAGE,
max_page_bytes: recall_wire::audit::MAX_PAGE_BYTES as u64,
})
.unwrap_or_default(),
);
capabilities.insert(
"scopes".to_string(),
serde_json::json!({ "kinds": ["project", "global", "machine"] }),
);
capabilities.insert(
"limits".to_string(),
serde_json::json!({
"max_body_bytes": super::MAX_BODY_BYTES,
"rate_limit": {
"max": state.cfg.rate_limit_max,
"window_seconds": state.cfg.rate_limit_window.as_secs(),
},
}),
);
json(
StatusCode::OK,
&Discovery {
protocol: Protocol {
current: PROTOCOL,
supported: vec![PROTOCOL],
},
server: ServerInfo {
version: discovery::version_for(channel, revision.as_deref()),
build: Build {
channel: channel.to_string(),
revision,
created: discovery::created().map(str::to_string),
},
},
min_client: MIN_CLIENT.to_string(),
auth: Auth {
methods: vec![
discovery::AUTH_BEARER.to_string(),
discovery::AUTH_DEVICE_SIG.to_string(),
],
},
capabilities,
},
)
}
pub(super) async fn handle_admin_stats(State(state): State<Arc<AppState>>) -> Response {
let (projects, totals) = match state.store.admin_stats() {
Ok(v) => v,
Err(e) => return internal(e),
};
let last_backup_at = state.read().last_backup_at.clone();
json(
StatusCode::OK,
&AdminStats {
projects,
totals,
git_commit: state.cfg.git_commit.clone(),
last_backup_at,
},
)
}
pub(super) async fn not_found() -> Response {
error(StatusCode::NOT_FOUND, "not found")
}