Skip to content
Merged
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
59 changes: 45 additions & 14 deletions addon/controllers/operations/scheduler/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,17 @@ import { inject as service } from '@ember/service';
import { action, computed } from '@ember/object';
import { isNone } from '@ember/utils';
import { isValid as isValidDate } from 'date-fns';
import { later } from '@ember/runloop';
import { task } from 'ember-concurrency';
import isObject from '@fleetbase/ember-core/utils/is-object';
import isJson from '@fleetbase/ember-core/utils/is-json';
import createFullCalendarEventFromOrder from '../../../utils/create-full-calendar-event-from-order';
import createFullCalendarEventFromScheduleItem from '../../../utils/create-full-calendar-event-from-schedule-item';
import toCalendarDate from '../../../utils/to-calendar-date';

// The ember-core socket service subscribes ~300 ms after listen() is called.
const SOCKET_SUBSCRIBE_SETTLE_MS = 500;

/**
* OperationsSchedulerIndexController
*
Expand Down Expand Up @@ -722,28 +726,55 @@ export default class OperationsSchedulerIndexController extends Controller {
// Real-Time Socket Subscriptions
// -------------------------------------------------------------------------

// Names of the socket channels this board opened, so teardown closes only those
// and leaves every other subscription (chat, notifications, ...) alone.
_socketChannelNames = new Set();

@action async subscribeToRealTimeUpdates() {
const orgId = this.currentUser?.companyId ?? this.currentUser?.company?.id;
if (!orgId) return;
await this.socket.listen(`company.${orgId}.orders`, (payload) => this._handleOrderSocketEvent(payload));
this.drivers.forEach(async (driver) => {
await this.socket.listen(`driver.${driver.id}`, (payload) => this._handleDriverSocketEvent(payload));
});
const drivers = this.drivers ?? [];
await Promise.all(
drivers.map((driver) => {
if (!driver?.id) return;
const channelName = `driver.${driver.id}`;
if (this._socketChannelNames.has(channelName)) return;
this._socketChannelNames.add(channelName);
return this.socket.listen(channelName, (payload) => this._handleDriverSocketEvent(payload));
})
);
}

@action unsubscribeFromRealTimeUpdates() {
if (this.socket && typeof this.socket.closeChannels === 'function') {
this.socket.closeChannels();
const names = this._socketChannelNames;
this._socketChannelNames = new Set();
if (!this.socket || names.size === 0) return;

if (typeof this.socket.closeChannel === 'function') {
names.forEach((name) => this.socket.closeChannel(name));
return;
}

const remaining = this._closeOpenedChannels(names);
if (remaining.size === 0) return;

// The socket service subscribes after a short delay, so a channel requested just before
// teardown may not exist yet: sweep once more after that delay.
later(this, () => this._closeOpenedChannels(remaining), SOCKET_SUBSCRIBE_SETTLE_MS);
}

_handleOrderSocketEvent({ data } = {}) {
if (!data?.id) return;
try {
this.store.pushPayload('order', { order: data });
} catch {
/* ignore */
/**
* Close the socket's channels with the given names, except any the board has opened again
* since. Returns the names that had no channel yet.
*/
_closeOpenedChannels(names) {
const remaining = new Set([...names].filter((name) => !this._socketChannelNames.has(name)));
const channels = Array.isArray(this.socket.channels) ? this.socket.channels : [];
for (const channel of channels) {
if (!names.has(channel?.name) || this._socketChannelNames.has(channel.name)) continue;
remaining.delete(channel.name);
channel.close();
}

return remaining;
}

_handleDriverSocketEvent({ event, data } = {}) {
Expand Down
2 changes: 1 addition & 1 deletion addon/routes/operations/scheduler/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,7 @@ export default class OperationsSchedulerIndexRoute extends Route {
}

resetController(controller) {
// Close all socket channels when the dispatcher navigates away
// Close the board's own socket channels when the dispatcher navigates away
// to prevent memory leaks and stale event handlers.
controller.unsubscribeFromRealTimeUpdates();
}
Expand Down
2 changes: 1 addition & 1 deletion composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
"barryvdh/laravel-dompdf": "^3.1",
"brick/geo": "0.7.2",
"cknow/laravel-money": "^7.1",
"fleetbase/core-api": ">=1.6.65",
"fleetbase/core-api": "^1.6.69",
"geocoder-php/google-maps-places-provider": "^1.4",
"giggsey/libphonenumber-for-php": "^8.13",
"league/geotools": "^1.1.0",
Expand Down
2 changes: 2 additions & 0 deletions server/src/Http/Controllers/Api/v1/DriverController.php
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
use Fleetbase\FleetOps\Support\GeofenceIntersectionService;
use Fleetbase\FleetOps\Support\OSRM;
use Fleetbase\FleetOps\Support\ProfileAccountManager;
use Fleetbase\FleetOps\Support\TrackingPublisher;
use Fleetbase\FleetOps\Support\Utils;
use Fleetbase\Http\Controllers\Controller;
use Fleetbase\Http\Requests\SwitchOrganizationRequest;
Expand Down Expand Up @@ -374,6 +375,7 @@ public function track(string $id, Request $request)
}

broadcast(new DriverLocationChanged($driver));
app(TrackingPublisher::class)->driverMoved($driver);

// ----------------------------------------------------------------
// Geofence intersection detection
Expand Down
2 changes: 2 additions & 0 deletions server/src/Http/Controllers/Api/v1/VehicleController.php
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
use Fleetbase\FleetOps\Models\Vendor;
use Fleetbase\FleetOps\Models\Warranty;
use Fleetbase\FleetOps\Support\GeofenceIntersectionService;
use Fleetbase\FleetOps\Support\TrackingPublisher;
use Fleetbase\FleetOps\Support\Utils;
use Fleetbase\Http\Controllers\Controller;
use Fleetbase\LaravelMysqlSpatial\Types\Point;
Expand Down Expand Up @@ -298,6 +299,7 @@ public function track(string $id, Request $request)
$vehicle->createPosition($positionData);

broadcast(new VehicleLocationChanged($vehicle));
app(TrackingPublisher::class)->vehicleMoved($vehicle);

try {
$newLocation = new Point($latitude, $longitude);
Expand Down
22 changes: 22 additions & 0 deletions server/src/Http/Controllers/Internal/v1/PositionController.php
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,10 @@
use Fleetbase\FleetOps\Jobs\ReplayPositions;
use Fleetbase\FleetOps\Models\Position;
use Fleetbase\FleetOps\Support\Utils;
use Fleetbase\Models\User;
use Fleetbase\Support\SocketCluster\ChannelAuthorizer;
use Fleetbase\Support\SocketCluster\SocketPrincipal;
use Fleetbase\Support\SocketCluster\SocketToken;
use Illuminate\Http\Request;

class PositionController extends FleetOpsController
Expand Down Expand Up @@ -37,6 +41,11 @@ public function replay(Request $request)
return response()->error('Position IDs are required');
}

// With socket authentication on, the replay may only publish to a channel the user could subscribe to.
if (SocketToken::enabled() && !$this->canReplayTo($request, $channelId)) {
return response()->error('You are not allowed to replay positions to this channel.', 403);
}

$positions = Position::whereIn('uuid', $positionIds)
->where('company_uuid', session('company'))
->orderBy('created_at')
Expand All @@ -57,6 +66,19 @@ public function replay(Request $request)
]);
}

/**
* Whether the requesting user may subscribe to, and so receive a replay on, the channel.
*/
protected function canReplayTo(Request $request, mixed $channelId): bool
{
$user = $request->user();
if (!$user instanceof User || !is_string($channelId) || $channelId === '' || strlen($channelId) > 255 || preg_match('/\s/', $channelId)) {
return false;
}

return app(ChannelAuthorizer::class)->authorize(SocketPrincipal::forUser($user, session('company')), $channelId)->allow;
}

/**
* Get position statistics/metrics.
*
Expand Down
50 changes: 50 additions & 0 deletions server/src/Jobs/PublishTrackingUpdate.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
<?php

namespace Fleetbase\FleetOps\Jobs;

use Fleetbase\FleetOps\Models\Order;
use Fleetbase\FleetOps\Support\TrackingPublisher;
use Fleetbase\FleetOps\Support\TrackingUpdate;
use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Bus\Dispatchable;
use Illuminate\Queue\InteractsWithQueue;
use Illuminate\Support\Facades\Cache;

/**
* Builds and publishes an order's public tracking updates off the request path.
*
* Location jobs publish the position only, except that at most once per ETA interval they
* publish the full update instead so that customers' ETAs follow the vehicle.
*/
class PublishTrackingUpdate implements ShouldQueue
{
use Dispatchable;
use InteractsWithQueue;
use Queueable;

public $tries = 1;
public $timeout = 60;

public function __construct(public string $orderUuid, public string $reason = TrackingUpdate::REASON_STATUS, public ?string $dbConnection = null)
{
}

public function handle(TrackingPublisher $publisher): int
{
$order = Order::on($this->dbConnection)->where('uuid', $this->orderUuid)->first();
if (!$order) {
return 0;
}

if ($this->reason !== TrackingUpdate::REASON_LOCATION) {
return $publisher->publishOrder($order, $this->reason);
}

if (Cache::add('fleetops:tracking:eta:' . $order->uuid, 1, TrackingPublisher::ETA_INTERVAL)) {
return $publisher->publishOrder($order, TrackingUpdate::REASON_ETA);
}

return $publisher->publishLocation($order);
}
}
27 changes: 27 additions & 0 deletions server/src/Listeners/PublishTrackingUpdates.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
<?php

namespace Fleetbase\FleetOps\Listeners;

use Fleetbase\FleetOps\Support\TrackingPublisher;

/**
* Queues public tracking updates when a waypoint's or an entity's activity changes.
*
* Runs synchronously but only looks up the order and queues a job, and does nothing at all
* while socket authentication is off.
*/
class PublishTrackingUpdates
{
public function __construct(protected ?TrackingPublisher $publisher = null)
{
$this->publisher ??= app(TrackingPublisher::class);
}

/**
* @param \Fleetbase\FleetOps\Events\WaypointActivityChanged|\Fleetbase\FleetOps\Events\WaypointCompleted|\Fleetbase\FleetOps\Events\EntityActivityChanged|\Fleetbase\FleetOps\Events\EntityCompleted $event
*/
public function handle(object $event): bool
{
return $this->publisher->stopChanged($event->waypoint ?? $event->entity ?? null);
}
}
4 changes: 4 additions & 0 deletions server/src/Observers/OrderObserver.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

use Fleetbase\FleetOps\Models\Order;
use Fleetbase\FleetOps\Support\LiveCacheService;
use Fleetbase\FleetOps\Support\TrackingPublisher;
use Illuminate\Support\Facades\Cache;

class OrderObserver
Expand Down Expand Up @@ -46,6 +47,9 @@ public function updated(Order $order)
}

$this->invalidateCache($order);

// Public tracking channels follow status, assignment and ETA changes.
app(TrackingPublisher::class)->orderChanged($order);
}

/**
Expand Down
11 changes: 11 additions & 0 deletions server/src/Providers/EventServiceProvider.php
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,17 @@ class EventServiceProvider extends ServiceProvider
\Fleetbase\FleetOps\Events\OrderFailed::class => [\Fleetbase\Listeners\SendResourceLifecycleWebhook::class, \Fleetbase\FleetOps\Listeners\NotifyOrderEvent::class],
\Fleetbase\FleetOps\Events\OrderReady::class => [\Fleetbase\FleetOps\Listeners\HandleOrderReady::class],

/*
* Stop Events
*
* Public tracking channels follow waypoint and entity activity. Order status, assignment
* and ETA changes are published from the order observer; positions from the tracking endpoints.
*/
\Fleetbase\FleetOps\Events\WaypointActivityChanged::class => [\Fleetbase\FleetOps\Listeners\PublishTrackingUpdates::class],
\Fleetbase\FleetOps\Events\WaypointCompleted::class => [\Fleetbase\FleetOps\Listeners\PublishTrackingUpdates::class],
\Fleetbase\FleetOps\Events\EntityActivityChanged::class => [\Fleetbase\FleetOps\Listeners\PublishTrackingUpdates::class],
\Fleetbase\FleetOps\Events\EntityCompleted::class => [\Fleetbase\FleetOps\Listeners\PublishTrackingUpdates::class],

/*
* Geofence Events
*
Expand Down
11 changes: 11 additions & 0 deletions server/src/Providers/FleetOpsServiceProvider.php
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,7 @@ public function boot()
});
$this->registerNotifications();
$this->registerAiCapabilities();
$this->registerSocketChannels();
$this->registerExpansionsFrom(__DIR__ . '/../Expansions');

// Register built-in orchestration engines.
Expand Down Expand Up @@ -265,6 +266,16 @@ public function registerNotifications()
]);
}

/**
* Register who may subscribe to FleetOps realtime channels, and the driver socket principal.
*/
protected function registerSocketChannels(): void
{
$this->callAfterResolving(\Fleetbase\Support\SocketCluster\SocketChannelRegistry::class, function (\Fleetbase\Support\SocketCluster\SocketChannelRegistry $registry) {
\Fleetbase\FleetOps\Support\SocketChannels::register($registry);
});
}

protected function registerAiCapabilities(): void
{
if (!Utils::classExists(\Fleetbase\Ai\Support\AiCapabilityRegistry::class)) {
Expand Down
Loading
Loading