Skip to main content

wire/
lib.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Shared protocol/auth transport types.
3//! Live authorization scope rules are owned by weft-server/src/access/scope.rs.
4
5#[cfg(test)]
6mod auth_tests;
7mod auth_token;
8mod key_binding;
9mod message_hosted;
10mod message_objects;
11mod message_pushpull;
12mod message_refs;
13mod message_status;
14mod native_pack;
15mod object_transfer;
16mod provider_pack;
17mod semantic_graph;
18
19pub use auth_token::AuthToken;
20pub use key_binding::{
21    WireKeyBinding, WireKeyBindingLiveness, WireKeyBindingRegistry, decode_key_binding_registry,
22    encode_key_binding_registry,
23};
24pub use message_hosted::{
25    HarnessIdentity, HostedGrantInfo, HostedNamespaceInfo, HostedRepositoryInfo, HostedSpoolInfo,
26    HostedSpoolKind, ProgressCheckpoint, SessionDiffSummary, SessionReportEnvelope,
27    TranscriptAttachmentRef, UsageTotals, WorktreeChangeBaseline,
28};
29pub use message_objects::{ObjectData, ObjectRequest};
30pub use message_pushpull::{PullComplete, PushComplete};
31pub use message_refs::{
32    AdvertisedRef, AdvertisedRefError, HeadInfo, RefEntry, RefFilter, RefKind, RefUpdated,
33};
34pub use message_status::{
35    Error, ErrorCode, RemoteCursorFailure, RemoteCursorReason, RemoteDuration, RemoteFailureCode,
36    RemoteFailureDetail, RemoteTimestamp,
37};
38pub use native_pack::{
39    GitPackChunkState, GrowingPackChunkReader, MAX_RECEIVED_GIT_PACK_SIZE,
40    MAX_RECEIVED_PACK_INDEX_SIZE, MAX_RECEIVED_PACK_SIZE, NativePackBundle, NativePackFileBundle,
41    NativePackStreamingWriter, PackChunkSpool, PackChunkState, PackFileChunkReader,
42    build_native_pack, install_received_pack, is_native_packable_object_type,
43    native_pack_excluded_object_types, next_pack_chunk, receive_pack_chunk,
44    reuse_native_pack_encoded_subset_in,
45};
46pub use object_transfer::{
47    MAX_PULL_FRAME_MESSAGE_SIZE, MAX_RECEIVED_REDACTIONS_BLOB_SIZE,
48    MAX_RECEIVED_STATE_VISIBILITY_BLOB_SIZE, admit_declared_received_len,
49    check_received_transfer_blob_size, chunk_bounds, chunk_count, chunk_offset, load_object_data,
50    load_requested_object, store_received_object,
51};
52pub use objects::transfer::{
53    GitLaneTransferIntent, ObjectAvailabilityPlan, ObjectId, ObjectInfo, ObjectType,
54    ObjectTypeBucket, PlannedObject, RepositoryTransferPlan, StateClosureOptions,
55    TransferPartitions, TransferPlanStats, enumerate_state_closure, enumerate_state_closure_plan,
56    enumerate_state_closure_plan_with_options, enumerate_state_closure_transfer_from_boundaries,
57    enumerate_state_closure_transfer_with_options, enumerate_state_closure_with_options,
58    has_object, is_ancestor, missing_blobs_in_tree, plan_object_availability,
59};
60pub use provider_pack::{
61    CompletedProviderPack, ProviderPackBundle, ProviderPackExtent, ProviderPackIndexEntry,
62    ProviderPackManifest, ProviderPackSpool, ProviderPackWriter, assemble_provider_pack,
63};
64pub use semantic_graph::{
65    SemanticGraphQueryKind, SemanticGraphQueryRequest, SemanticGraphQueryResponse, SemanticGraphRef,
66};
67
68/// Default port for Heddle protocol.
69pub const DEFAULT_PORT: u16 = 8421;
70
71/// Protocol version.
72pub const PROTOCOL_VERSION: u32 = 1;
73
74/// Maximum message size (64 MB).
75pub const MAX_MESSAGE_SIZE: usize = 64 * 1024 * 1024;
76
77/// Error type for protocol operations.
78#[derive(Debug, thiserror::Error)]
79pub enum ProtocolError {
80    #[error("foreign prefix {limit_name} exceeds limit {limit}")]
81    ForeignPrefixLimitExceeded {
82        limit_name: &'static str,
83        limit: usize,
84    },
85    #[error("io error: {0}")]
86    Io(#[from] std::io::Error),
87
88    #[error("serialization error: {0}")]
89    Serialization(String),
90
91    #[error("message too large: {size} bytes (max {max})")]
92    MessageTooLarge { size: usize, max: usize },
93
94    #[error("invalid message type: {0}")]
95    InvalidMessageType(u8),
96
97    #[error("protocol version mismatch: server={server}, client={client}")]
98    VersionMismatch { server: u32, client: u32 },
99
100    #[error("capability not supported: {0}")]
101    CapabilityNotSupported(String),
102
103    #[error("authentication failed: {0}")]
104    AuthenticationFailed(String),
105
106    #[error("authorization failed: {0}")]
107    AuthorizationFailed(String),
108
109    #[error("object not found: {0}")]
110    ObjectNotFound(String),
111
112    #[error("already exists: {0}")]
113    AlreadyExists(String),
114
115    #[error("invalid state: {0}")]
116    InvalidState(String),
117
118    #[error(
119        "Thread '{thread}' has no capture of its own; run `heddle capture -m \"...\"` in its checkout, then push again"
120    )]
121    ThreadCaptureRequired { thread: String },
122
123    #[error(
124        "Thread {thread} has {states} States and {operations} operations; {limit_name} limit is {limit}, required {actual}. Very large histories need server-side incremental admission (HeddleCo/weft#2432) or hosted import"
125    )]
126    PublicationLimitExceeded {
127        thread: String,
128        states: usize,
129        operations: usize,
130        limit_name: &'static str,
131        limit: usize,
132        actual: usize,
133    },
134
135    #[error(
136        "Thread {thread} operation {operation} requires {bytes} encoded bytes; one publication batch allows {batch_limit} bytes and the negotiated frame allows {frame_limit} bytes"
137    )]
138    PublicationOperationTooLarge {
139        thread: String,
140        operation: usize,
141        bytes: usize,
142        batch_limit: usize,
143        frame_limit: usize,
144    },
145
146    #[error("remote error: {0}")]
147    Remote(String),
148
149    #[error("remote failure ({code:?}): {message}")]
150    RemoteFailure {
151        code: RemoteFailureCode,
152        message: String,
153        details: Vec<RemoteFailureDetail>,
154    },
155
156    #[error("lock error: {0}")]
157    LockError(String),
158}
159
160impl From<rmp_serde::encode::Error> for ProtocolError {
161    fn from(e: rmp_serde::encode::Error) -> Self {
162        ProtocolError::Serialization(e.to_string())
163    }
164}
165
166impl From<rmp_serde::decode::Error> for ProtocolError {
167    fn from(e: rmp_serde::decode::Error) -> Self {
168        ProtocolError::Serialization(e.to_string())
169    }
170}
171
172impl From<objects::error::HeddleError> for ProtocolError {
173    fn from(e: objects::error::HeddleError) -> Self {
174        match &e {
175            // Missing objects must surface as NotFound, not degrade to a
176            // generic server error via the blanket Remote arm below.
177            objects::error::HeddleError::NotFound(_)
178            | objects::error::HeddleError::StateNotFound(_)
179            | objects::error::HeddleError::MissingObject { .. } => {
180                ProtocolError::ObjectNotFound(e.to_string())
181            }
182            _ => ProtocolError::Remote(e.to_string()),
183        }
184    }
185}
186
187impl ProtocolError {
188    pub fn client_message(&self) -> String {
189        match self {
190            ProtocolError::ForeignPrefixLimitExceeded { .. } => self.to_string(),
191            ProtocolError::Io(_) => "network error".to_string(),
192            ProtocolError::Serialization(_) => "protocol error".to_string(),
193            ProtocolError::MessageTooLarge { .. } => "message too large".to_string(),
194            ProtocolError::InvalidMessageType(_) => "protocol error".to_string(),
195            ProtocolError::VersionMismatch { .. } => "protocol version mismatch".to_string(),
196            ProtocolError::CapabilityNotSupported(_) => "capability not supported".to_string(),
197            ProtocolError::AuthenticationFailed(_) => "permission denied".to_string(),
198            ProtocolError::AuthorizationFailed(_) => "permission denied".to_string(),
199            ProtocolError::ObjectNotFound(_) => "object not found".to_string(),
200            ProtocolError::AlreadyExists(_) => "resource already exists".to_string(),
201            ProtocolError::InvalidState(_) => "invalid request state".to_string(),
202            ProtocolError::ThreadCaptureRequired { .. } => self.to_string(),
203            ProtocolError::PublicationLimitExceeded { .. }
204            | ProtocolError::PublicationOperationTooLarge { .. } => self.to_string(),
205            ProtocolError::Remote(_) => "internal server error".to_string(),
206            ProtocolError::RemoteFailure { message, .. } => message.clone(),
207            ProtocolError::LockError(_) => "internal server error".to_string(),
208        }
209    }
210
211    pub fn error_code(&self) -> ErrorCode {
212        match self {
213            ProtocolError::ForeignPrefixLimitExceeded { .. } => ErrorCode::InvalidArgument,
214            ProtocolError::Io(_) => ErrorCode::Network,
215            ProtocolError::Serialization(_) => ErrorCode::Protocol,
216            ProtocolError::MessageTooLarge { .. } => ErrorCode::Protocol,
217            ProtocolError::InvalidMessageType(_) => ErrorCode::Protocol,
218            ProtocolError::VersionMismatch { .. } => ErrorCode::Protocol,
219            ProtocolError::CapabilityNotSupported(_) => ErrorCode::Protocol,
220            ProtocolError::AuthenticationFailed(_) => ErrorCode::PermissionDenied,
221            ProtocolError::AuthorizationFailed(_) => ErrorCode::PermissionDenied,
222            ProtocolError::ObjectNotFound(_) => ErrorCode::NotFound,
223            ProtocolError::AlreadyExists(_) => ErrorCode::InvalidArgument,
224            ProtocolError::InvalidState(_) => ErrorCode::InvalidArgument,
225            ProtocolError::ThreadCaptureRequired { .. } => ErrorCode::InvalidArgument,
226            ProtocolError::PublicationLimitExceeded { .. }
227            | ProtocolError::PublicationOperationTooLarge { .. } => ErrorCode::InvalidArgument,
228            ProtocolError::Remote(_) => ErrorCode::Server,
229            ProtocolError::RemoteFailure { code, .. } => match code {
230                RemoteFailureCode::InvalidArgument
231                | RemoteFailureCode::AlreadyExists
232                | RemoteFailureCode::FailedPrecondition
233                | RemoteFailureCode::OutOfRange => ErrorCode::InvalidArgument,
234                RemoteFailureCode::NotFound => ErrorCode::NotFound,
235                RemoteFailureCode::PermissionDenied | RemoteFailureCode::Unauthenticated => {
236                    ErrorCode::PermissionDenied
237                }
238                RemoteFailureCode::DeadlineExceeded
239                | RemoteFailureCode::ResourceExhausted
240                | RemoteFailureCode::Aborted
241                | RemoteFailureCode::Unavailable
242                | RemoteFailureCode::Cancelled => ErrorCode::Network,
243                RemoteFailureCode::Unspecified
244                | RemoteFailureCode::Unknown
245                | RemoteFailureCode::Unimplemented
246                | RemoteFailureCode::Internal
247                | RemoteFailureCode::DataLoss => ErrorCode::Server,
248            },
249            ProtocolError::LockError(_) => ErrorCode::Server,
250        }
251    }
252
253    pub fn to_wire_error(&self, details: Option<String>) -> Error {
254        Error {
255            code: self.error_code(),
256            message: self.client_message(),
257            details,
258        }
259    }
260}
261
262pub type Result<T> = std::result::Result<T, ProtocolError>;
263
264#[cfg(test)]
265mod tests {
266    use std::io;
267
268    use super::{ErrorCode, ProtocolError, RemoteFailureCode};
269
270    #[test]
271    fn missing_object_errors_surface_as_not_found() {
272        use objects::error::HeddleError;
273
274        let cases = vec![
275            HeddleError::NotFound("blob missing".to_string()),
276            HeddleError::StateNotFound(objects::object::StateId::from_content_hash(
277                objects::object::ContentHash::compute(b"missing"),
278            )),
279            HeddleError::MissingObject {
280                object_type: "tree".to_string(),
281                id: "abc123".to_string(),
282            },
283        ];
284        for err in cases {
285            let protocol_error = ProtocolError::from(err);
286            assert_eq!(
287                protocol_error.error_code(),
288                ErrorCode::NotFound,
289                "missing-object errors must surface as NotFound, got {protocol_error:?}"
290            );
291            assert!(
292                !matches!(protocol_error, ProtocolError::Remote(_)),
293                "missing-object errors must not degrade to Remote/Server"
294            );
295        }
296    }
297
298    #[test]
299    fn protocol_error_public_mapping_is_stable() {
300        let cases = vec![
301            (
302                ProtocolError::Io(io::Error::new(io::ErrorKind::TimedOut, "timeout")),
303                "network error",
304                ErrorCode::Network,
305            ),
306            (
307                ProtocolError::Serialization("bad msgpack".to_string()),
308                "protocol error",
309                ErrorCode::Protocol,
310            ),
311            (
312                ProtocolError::MessageTooLarge { size: 65, max: 64 },
313                "message too large",
314                ErrorCode::Protocol,
315            ),
316            (
317                ProtocolError::InvalidMessageType(42),
318                "protocol error",
319                ErrorCode::Protocol,
320            ),
321            (
322                ProtocolError::VersionMismatch {
323                    server: 2,
324                    client: 1,
325                },
326                "protocol version mismatch",
327                ErrorCode::Protocol,
328            ),
329            (
330                ProtocolError::CapabilityNotSupported("pack-v2".to_string()),
331                "capability not supported",
332                ErrorCode::Protocol,
333            ),
334            (
335                ProtocolError::AuthenticationFailed("bad token".to_string()),
336                "permission denied",
337                ErrorCode::PermissionDenied,
338            ),
339            (
340                ProtocolError::AuthorizationFailed("missing grant".to_string()),
341                "permission denied",
342                ErrorCode::PermissionDenied,
343            ),
344            (
345                ProtocolError::ObjectNotFound("abc123".to_string()),
346                "object not found",
347                ErrorCode::NotFound,
348            ),
349            (
350                ProtocolError::AlreadyExists("__users/luke/repo".to_string()),
351                "resource already exists",
352                ErrorCode::InvalidArgument,
353            ),
354            (
355                ProtocolError::InvalidState("bad resume".to_string()),
356                "invalid request state",
357                ErrorCode::InvalidArgument,
358            ),
359            (
360                ProtocolError::Remote("database unavailable".to_string()),
361                "internal server error",
362                ErrorCode::Server,
363            ),
364            (
365                ProtocolError::RemoteFailure {
366                    code: RemoteFailureCode::InvalidArgument,
367                    message: "server supplied message".to_string(),
368                    details: Vec::new(),
369                },
370                "server supplied message",
371                ErrorCode::InvalidArgument,
372            ),
373            (
374                ProtocolError::LockError("ref locked".to_string()),
375                "internal server error",
376                ErrorCode::Server,
377            ),
378        ];
379
380        for (error, expected_message, expected_code) in cases {
381            assert_eq!(error.client_message(), expected_message);
382            assert_eq!(error.error_code(), expected_code);
383
384            let wire_error = error.to_wire_error(Some("trace id".to_string()));
385            assert_eq!(wire_error.code, expected_code);
386            assert_eq!(wire_error.message, expected_message);
387            assert_eq!(wire_error.details.as_deref(), Some("trace id"));
388        }
389    }
390}