Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions bcmp/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,9 @@ set(SOURCES
dfu_client.c
dfu_core.c
dfu_host.c
ftp_core.c
ftp_coordinator.c
ftp_endpoint.c
heartbeat.c
info.c
neighbors.c
Expand Down
2 changes: 2 additions & 0 deletions bcmp/bcmp.c
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
#include "bm_ip.h"
#include "bm_os.h"
#include "dfu.h"
#include "ftp.h"
#include "l2.h"
#include "messages/config.h"
#include "messages/heartbeat.h"
Expand Down Expand Up @@ -122,6 +123,7 @@ BmErr bcmp_init(NetworkDevice network_device) {
bm_err_check(err, ping_init());
bm_err_check(err, time_init());
bm_err_check(err, bm_dfu_init());
bm_err_check(err, bm_ftp_init());
bm_err_check(err, bcmp_config_init());
bm_err_check(err, bcmp_neighbor_init(network_device.trait->num_ports()));
bm_err_check(err, bcmp_device_info_init());
Expand Down
65 changes: 65 additions & 0 deletions bcmp/ftp.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
#pragma once

#ifdef __cplusplus
extern "C" {
#endif

#include "bm_os.h"
#include "messages.h"

#ifndef BM_FTP_EVENT_TASK_PRIORITY
#define BM_FTP_EVENT_TASK_PRIORITY 11
#endif

typedef enum {
BmFtpEventStart = 0,
BmFtpEventFetch,
BmFtpEventAck,
BmFtpEventChunkRequest,
BmFtpEventChunk,
BmFtpEventEnd,
BmFtpEventAbort,
} BmFtpEventType;

typedef struct {
BmFtpEventType type;
uint8_t *payload;
size_t payload_len;
} BmFtpEvent;

typedef BmErr (*BmFtpEventHandler)(const BmFtpEvent *event, void *context);

BmErr bm_ftp_init(void);
BmErr bm_ftp_coordinator_init(void);
BmErr bm_ftp_coordinator_process_event(const BmFtpEvent *event, void *context);
BmQueue bm_ftp_get_event_queue(void);
void bm_ftp_set_event_handler(BmFtpEventHandler handler, void *context);

BmErr bm_ftp_start_fetch(uint64_t source_node_id, uint32_t transfer_id,
BmFtpEndpointKind source_kind,
const uint8_t *source_spec, uint16_t source_spec_len,
BmFtpEndpointKind sink_kind, const uint8_t *sink_spec,
uint16_t sink_spec_len);
BmErr bm_ftp_send_ack(uint64_t dst_node_id, uint32_t transfer_id, bool success,
BmFtpErr error, uint32_t total_size, uint16_t crc16,
uint16_t chunk_size);
BmErr bm_ftp_send_start(uint64_t dst_node_id, uint32_t transfer_id,
uint32_t total_size, uint16_t requested_chunk_size,
uint16_t crc16, BmFtpEndpointKind sink_kind,
const uint8_t *sink_spec, uint16_t sink_spec_len);
BmErr bm_ftp_send_end(uint64_t dst_node_id, uint32_t transfer_id, bool success,
BmFtpErr error, uint32_t bytes_received, uint16_t running_crc16);
BmErr bm_ftp_send_fetch(uint64_t dst_node_id, uint32_t transfer_id,
BmFtpEndpointKind source_kind,
const uint8_t *source_spec, uint16_t source_spec_len);
BmErr bm_ftp_request_chunk(uint64_t dst_node_id, uint32_t transfer_id,
uint32_t offset, uint16_t length);
BmErr bm_ftp_send_chunk(uint64_t dst_node_id, uint32_t transfer_id,
uint32_t offset, const uint8_t *payload,
uint16_t payload_len);
BmErr bm_ftp_send_abort(uint64_t dst_node_id, uint32_t transfer_id,
BmFtpErr error);

#ifdef __cplusplus
}
#endif
300 changes: 300 additions & 0 deletions bcmp/ftp_coordinator.c
Original file line number Diff line number Diff line change
@@ -0,0 +1,300 @@
#include "ftp.h"

#include "device.h"
#include "ftp_endpoint.h"
#include <string.h>

#define bm_ftp_max_chunk_size (1024U)
#define bm_ftp_max_endpoint_spec_len (128U)

typedef enum {
BmFtpCoordinatorIdle = 0,
BmFtpCoordinatorWaitingAck,
BmFtpCoordinatorServing,
BmFtpCoordinatorReceiving,
} BmFtpCoordinatorState;

typedef struct {
BmFtpCoordinatorState state;
uint64_t peer_node_id;
uint32_t transfer_id;
uint32_t total_size;
uint32_t bytes_received;
uint16_t chunk_size;
uint16_t crc16;
BmFtpEndpointKind pending_sink_kind;
uint16_t pending_sink_spec_len;
uint8_t pending_sink_spec[bm_ftp_max_endpoint_spec_len];
BmFtpSource source;
BmFtpSink sink;
} BmFtpCoordinator;

static BmFtpCoordinator coordinator;

static BmFtpErr bm_ftp_error_from_bm_error(BmErr error) {
if (error == BmENODEV) {
return BmFtpErrUnsupported;
}
if (error == BmEINVAL) {
return BmFtpErrInvalidSpec;
}
return BmFtpErrWriteFailed;
}

static void bm_ftp_reset(void) {
if (coordinator.state == BmFtpCoordinatorServing) {
bm_ftp_source_close(&coordinator.source);
} else if (coordinator.state == BmFtpCoordinatorReceiving) {
bm_ftp_sink_close(&coordinator.sink);
}
memset(&coordinator, 0, sizeof(coordinator));
}

static BmErr bm_ftp_reject(const BmFtpAddress *address, uint32_t transfer_id,
BmFtpErr error) {
return bm_ftp_send_ack(address->src_node_id, transfer_id, false, error, 0, 0, 0);
}

static BmErr bm_ftp_handle_start(const BmFtpStart *start) {
if (coordinator.state != BmFtpCoordinatorIdle) {
return bm_ftp_reject(&start->addresses, start->transfer_id, BmFtpErrBusy);
}
if (start->total_size == 0 || start->requested_chunk_size == 0 ||
start->requested_chunk_size > bm_ftp_max_chunk_size) {
return bm_ftp_reject(&start->addresses, start->transfer_id, BmFtpErrTooLarge);
}

BmErr error = bm_ftp_sink_open((BmFtpEndpointKind)start->sink_kind, start->sink_spec,
start->sink_spec_len, start->total_size, &coordinator.sink);
if (error != BmOK) {
return bm_ftp_reject(&start->addresses, start->transfer_id,
bm_ftp_error_from_bm_error(error));
}

coordinator.state = BmFtpCoordinatorReceiving;
coordinator.peer_node_id = start->addresses.src_node_id;
coordinator.transfer_id = start->transfer_id;
coordinator.total_size = start->total_size;
coordinator.chunk_size = start->requested_chunk_size;
coordinator.crc16 = start->crc16;

error = bm_ftp_send_ack(coordinator.peer_node_id, coordinator.transfer_id, true,
BmFtpErrNone, coordinator.total_size, coordinator.crc16,
coordinator.chunk_size);
if (error != BmOK) {
bm_ftp_reset();
return error;
}
uint16_t requested_length = coordinator.total_size < coordinator.chunk_size
? (uint16_t)coordinator.total_size
: coordinator.chunk_size;
return bm_ftp_request_chunk(coordinator.peer_node_id, coordinator.transfer_id, 0,
requested_length);
}

static BmErr bm_ftp_handle_fetch(const BmFtpFetch *fetch) {
if (coordinator.state != BmFtpCoordinatorIdle) {
return bm_ftp_reject(&fetch->addresses, fetch->transfer_id, BmFtpErrBusy);
}

BmErr error = bm_ftp_source_open((BmFtpEndpointKind)fetch->source_kind,
fetch->source_spec, fetch->source_spec_len,
&coordinator.source);
if (error != BmOK) {
return bm_ftp_reject(&fetch->addresses, fetch->transfer_id,
bm_ftp_error_from_bm_error(error));
}

coordinator.state = BmFtpCoordinatorServing;
coordinator.peer_node_id = fetch->addresses.src_node_id;
coordinator.transfer_id = fetch->transfer_id;
coordinator.total_size = coordinator.source.total_size;
coordinator.crc16 = coordinator.source.crc16;
coordinator.chunk_size = bm_ftp_max_chunk_size;

if (coordinator.total_size == 0) {
bm_ftp_reset();
return bm_ftp_reject(&fetch->addresses, fetch->transfer_id, BmFtpErrInvalidSpec);
}

error = bm_ftp_send_ack(coordinator.peer_node_id, coordinator.transfer_id, true,
BmFtpErrNone, coordinator.total_size, coordinator.crc16,
coordinator.chunk_size);
if (error != BmOK) {
bm_ftp_reset();
}
return error;
}

static BmErr bm_ftp_handle_ack(const BmFtpAck *ack) {
if (coordinator.state != BmFtpCoordinatorWaitingAck ||
ack->addresses.src_node_id != coordinator.peer_node_id ||
ack->transfer_id != coordinator.transfer_id) {
return BmEBADMSG;
}
if (!ack->success) {
bm_ftp_reset();
return BmECONNREFUSED;
}
if (ack->total_size == 0 || ack->chunk_size == 0 ||
ack->chunk_size > bm_ftp_max_chunk_size) {
bm_ftp_reset();
return BmEBADMSG;
}

BmErr error = bm_ftp_sink_open(coordinator.pending_sink_kind,
coordinator.pending_sink_spec,
coordinator.pending_sink_spec_len,
ack->total_size, &coordinator.sink);
if (error != BmOK) {
bm_ftp_send_abort(coordinator.peer_node_id, coordinator.transfer_id,
bm_ftp_error_from_bm_error(error));
bm_ftp_reset();
return error;
}
coordinator.state = BmFtpCoordinatorReceiving;
coordinator.total_size = ack->total_size;
coordinator.crc16 = ack->crc16;
coordinator.chunk_size = ack->chunk_size;
uint16_t requested_length = coordinator.total_size < coordinator.chunk_size
? (uint16_t)coordinator.total_size
: coordinator.chunk_size;
return bm_ftp_request_chunk(coordinator.peer_node_id, coordinator.transfer_id, 0,
requested_length);
}

static BmErr bm_ftp_handle_chunk_request(const BmFtpChunkRequest *request) {
if (coordinator.state != BmFtpCoordinatorServing ||
request->addresses.src_node_id != coordinator.peer_node_id ||
request->transfer_id != coordinator.transfer_id || request->length == 0 ||
request->length > coordinator.chunk_size || request->offset >= coordinator.total_size ||
request->length > coordinator.total_size - request->offset) {
return BmEBADMSG;
}

uint8_t buffer[bm_ftp_max_chunk_size];
BmErr error = bm_ftp_source_read_at(&coordinator.source, request->offset, buffer,
request->length);
if (error != BmOK) {
bm_ftp_send_abort(coordinator.peer_node_id, coordinator.transfer_id,
bm_ftp_error_from_bm_error(error));
bm_ftp_reset();
return error;
}
return bm_ftp_send_chunk(coordinator.peer_node_id, coordinator.transfer_id, request->offset,
buffer, request->length);
}

static BmErr bm_ftp_handle_chunk(const BmFtpChunk *chunk) {
if (coordinator.state != BmFtpCoordinatorReceiving ||
chunk->addresses.src_node_id != coordinator.peer_node_id ||
chunk->transfer_id != coordinator.transfer_id ||
chunk->offset != coordinator.bytes_received || chunk->payload_length == 0 ||
chunk->payload_length > coordinator.chunk_size ||
chunk->payload_length > coordinator.total_size - coordinator.bytes_received) {
return BmEBADMSG;
}

BmErr error = bm_ftp_sink_write_at(&coordinator.sink, chunk->offset, chunk->payload,
chunk->payload_length);
if (error != BmOK) {
bm_ftp_send_abort(coordinator.peer_node_id, coordinator.transfer_id,
bm_ftp_error_from_bm_error(error));
bm_ftp_sink_abort(&coordinator.sink);
bm_ftp_reset();
return error;
}
coordinator.bytes_received += chunk->payload_length;
if (coordinator.bytes_received == coordinator.total_size) {
error = bm_ftp_sink_finalize(&coordinator.sink, coordinator.total_size, coordinator.crc16);
BmFtpErr ftp_error = error == BmOK ? BmFtpErrNone : bm_ftp_error_from_bm_error(error);
bm_ftp_send_end(coordinator.peer_node_id, coordinator.transfer_id, error == BmOK, ftp_error,
coordinator.bytes_received, coordinator.crc16);
bm_ftp_reset();
return error;
}

uint32_t remaining = coordinator.total_size - coordinator.bytes_received;
uint16_t requested_length = remaining < coordinator.chunk_size ? (uint16_t)remaining
: coordinator.chunk_size;
return bm_ftp_request_chunk(coordinator.peer_node_id, coordinator.transfer_id,
coordinator.bytes_received, requested_length);
}

static BmErr bm_ftp_handle_end(const BmFtpEnd *end) {
if (coordinator.state != BmFtpCoordinatorServing ||
end->addresses.src_node_id != coordinator.peer_node_id ||
end->transfer_id != coordinator.transfer_id) {
return BmEBADMSG;
}
bm_ftp_reset();
return end->success ? BmOK : BmECANCELED;
}

static BmErr bm_ftp_handle_abort(const BmFtpAbort *abort) {
if (coordinator.state == BmFtpCoordinatorIdle ||
abort->addresses.src_node_id != coordinator.peer_node_id ||
abort->transfer_id != coordinator.transfer_id) {
return BmEBADMSG;
}
if (coordinator.state == BmFtpCoordinatorReceiving) {
bm_ftp_sink_abort(&coordinator.sink);
}
bm_ftp_reset();
return BmECANCELED;
}

BmErr bm_ftp_coordinator_process_event(const BmFtpEvent *event, void *context) {
(void)context;
switch (event->type) {
case BmFtpEventStart:
return bm_ftp_handle_start((const BmFtpStart *)event->payload);
case BmFtpEventFetch:
return bm_ftp_handle_fetch((const BmFtpFetch *)event->payload);
case BmFtpEventAck:
return bm_ftp_handle_ack((const BmFtpAck *)event->payload);
case BmFtpEventChunkRequest:
return bm_ftp_handle_chunk_request((const BmFtpChunkRequest *)event->payload);
case BmFtpEventChunk:
return bm_ftp_handle_chunk((const BmFtpChunk *)event->payload);
case BmFtpEventEnd:
return bm_ftp_handle_end((const BmFtpEnd *)event->payload);
case BmFtpEventAbort:
return bm_ftp_handle_abort((const BmFtpAbort *)event->payload);
default:
return BmOK;
}
}

BmErr bm_ftp_coordinator_init(void) {
memset(&coordinator, 0, sizeof(coordinator));
bm_ftp_set_event_handler(bm_ftp_coordinator_process_event, NULL);
return BmOK;
}

BmErr bm_ftp_start_fetch(uint64_t source_node_id, uint32_t transfer_id,
BmFtpEndpointKind source_kind, const uint8_t *source_spec,
uint16_t source_spec_len, BmFtpEndpointKind sink_kind,
const uint8_t *sink_spec, uint16_t sink_spec_len) {
if (coordinator.state != BmFtpCoordinatorIdle ||
(source_spec_len && !source_spec) || (sink_spec_len && !sink_spec) ||
sink_spec_len > bm_ftp_max_endpoint_spec_len) {
return BmEINVAL;
}

coordinator.state = BmFtpCoordinatorWaitingAck;
coordinator.peer_node_id = source_node_id;
coordinator.transfer_id = transfer_id;
coordinator.pending_sink_kind = sink_kind;
coordinator.pending_sink_spec_len = sink_spec_len;
if (sink_spec_len) {
memcpy(coordinator.pending_sink_spec, sink_spec, sink_spec_len);
}

BmErr error = bm_ftp_send_fetch(source_node_id, transfer_id, source_kind, source_spec,
source_spec_len);
if (error != BmOK) {
bm_ftp_reset();
}
return error;
}
Loading
Loading